From eda176a82c547297eb7033942909febb6b55208b Mon Sep 17 00:00:00 2001 From: Michael Spilman <130705596+TheDrOnos@users.noreply.github.com> Date: Fri, 25 Sep 2026 08:56:57 -0500 Subject: [PATCH 01/13] framing: encoders for the PLOTAPP/PLOTBIN wire format, the inverse of the demuxer --- de_shell/js/main/framing.test.ts | 125 +++++++++++++++++++++++++++++++ de_shell/js/main/framing.ts | 34 +++++++++ 2 files changed, 159 insertions(+) create mode 100644 de_shell/js/main/framing.test.ts create mode 100644 de_shell/js/main/framing.ts diff --git a/de_shell/js/main/framing.test.ts b/de_shell/js/main/framing.test.ts new file mode 100644 index 0000000..386bc97 --- /dev/null +++ b/de_shell/js/main/framing.test.ts @@ -0,0 +1,125 @@ +/** + * framing.test.ts — the encoders against the existing demuxer: whatever + * framing.ts writes, stdoutDemux.ts must read back unchanged, whole or split + * one byte at a time. + * + * Run: `node --test de_shell/js/main/framing.test.ts`, or via `npm run test:unit`. + */ +import { test } from 'node:test' +import assert from 'node:assert/strict' +import { encodeBinary, encodeMessage } from './framing.ts' +import { createStdoutDemux } from './stdoutDemux.ts' + +type Event = + | { kind: 'message'; msg: Record } + | { kind: 'stream'; text: string } + | { kind: 'binary'; header: Record; payload: string } + +/** Feed `stream` to a fresh demuxer in `step`-byte slices; return the event trace. */ +function decode(stream: Buffer, step = stream.length || 1): Event[] { + const events: Event[] = [] + const demux = createStdoutDemux({ + onMessage: (msg) => events.push({ kind: 'message', msg }), + onStream: (text) => events.push({ kind: 'stream', text }), + onBinary: (header, payload) => + events.push({ kind: 'binary', header, payload: payload.toString('hex') }), + }) + for (let pos = 0; pos < stream.length; pos += step) { + demux.push(stream.subarray(pos, pos + step)) + } + return events +} + +/** A whole binary frame as it goes on the wire: prefix, header, then the caller's payload. */ +function frame(header: Record, payload: Buffer): Buffer { + return Buffer.concat([...encodeBinary(header, payload), payload]) +} + +function patternPayload(n: number): Buffer { + const p = Buffer.allocUnsafe(n) + for (let i = 0; i < n; i++) p[i] = i & 0xff // includes 0x0a bytes + return p +} + +test('a message is one PLOTAPP line and decodes to the same object', () => { + const msg = { type: 'state_update', key: 'clim', value: [0, 255], nested: { ok: true, none: null } } + const bytes = encodeMessage(msg) + assert.equal(bytes.subarray(0, 8).toString('ascii'), 'PLOTAPP:') + assert.equal(bytes.indexOf(0x0a), bytes.length - 1, 'exactly one newline, at the end') + assert.deepEqual(decode(bytes), [{ kind: 'message', msg }]) +}) + +test('a newline inside a string stays escaped, so the message is still one line', () => { + const msg = { type: 'status', text: 'line one\nline two\r\n' } + const bytes = encodeMessage(msg) + assert.equal(bytes.indexOf(0x0a), bytes.length - 1) + assert.deepEqual(decode(bytes), [{ kind: 'message', msg }]) +}) + +test('a frame with an empty payload decodes with an empty payload', () => { + assert.deepEqual(decode(frame({ fig_id: 'f3', key: 'spec' }, Buffer.alloc(0))), [ + { kind: 'binary', header: { fig_id: 'f3', key: 'spec' }, payload: '' }, + ]) +}) + +test('a payload holding newlines and PLOTBIN:/PLOTAPP: lookalikes decodes whole, at any chunking', () => { + const nasty = Buffer.concat([ + Buffer.from('\nPLOTBIN:9:9\nPLOTAPP:{}\n', 'ascii'), + patternPayload(3000), + ]) + const stream = Buffer.concat([ + frame({ fig_id: 'f1', key: 'image' }, nasty), + encodeMessage({ type: 'done' }), + ]) + const expected: Event[] = [ + { kind: 'binary', header: { fig_id: 'f1', key: 'image' }, payload: nasty.toString('hex') }, + { kind: 'message', msg: { type: 'done' } }, + ] + for (const step of [stream.length, 65536, 7, 1]) { + assert.deepEqual(decode(stream, step), expected, `diverged at ${step}-byte chunks`) + } +}) + +test('a non-ASCII header: hlen counts UTF-8 bytes, not characters', () => { + const header = { fig_id: 'f', key: 'k', label: 'εxx Å' } + const [prefix, head] = encodeBinary(header, Buffer.from([1, 2, 3])) + const json = JSON.stringify(header) + assert.notEqual(Buffer.byteLength(json, 'utf8'), json.length, 'the fixture must be multi-byte') + assert.equal(prefix.toString('ascii'), `PLOTBIN:${Buffer.byteLength(json, 'utf8')}:3\n`) + assert.equal(head.toString('utf8'), json) + assert.deepEqual(decode(frame(header, Buffer.from([1, 2, 3])), 1), [ + { kind: 'binary', header, payload: '010203' }, + ]) +}) + +test('encodeBinary returns the prefix and header only; the payload is never copied', () => { + const payload = patternPayload(64) + const parts = encodeBinary({ fig_id: 'f', key: 'k' }, payload) + assert.equal(parts.length, 2) + assert.ok(!parts.includes(payload)) + assert.match(parts[0].toString('ascii'), /^PLOTBIN:\d+:64\n$/) +}) + +test('the layout is the documented wire format, byte for byte', () => { + // PLOTBIN::\n
, as anyplotlib's + // _binary_frame.encode_frame writes it and stdoutDemux.ts parses it. + const header = { fig_id: 'f', key: 'k', dims: [2, 2] } + const payload = Buffer.from([0, 10, 255, 10]) + const json = JSON.stringify(header) + assert.deepEqual(frame(header, payload), Buffer.concat([ + Buffer.from(`PLOTBIN:${Buffer.byteLength(json)}:4\n${json}`, 'utf8'), + payload, + ])) + assert.deepEqual(encodeMessage({ type: 'done' }), Buffer.from('PLOTAPP:{"type":"done"}\n', 'utf8')) +}) + +test('a non-finite number is refused rather than framed', () => { + for (const bad of [NaN, Infinity, -Infinity]) { + assert.throws(() => encodeMessage({ type: 'fit', value: bad }), RangeError) + assert.throws( + () => encodeBinary({ fig_id: 'f', key: 'k', clim: [0, bad] }, Buffer.alloc(0)), + RangeError, + ) + } + assert.throws(() => encodeMessage({ a: { b: [1, NaN] } }), /non-finite/) +}) diff --git a/de_shell/js/main/framing.ts b/de_shell/js/main/framing.ts new file mode 100644 index 0000000..47c9ba9 --- /dev/null +++ b/de_shell/js/main/framing.ts @@ -0,0 +1,34 @@ +/** + * framing.ts — encoders for the backend wire format; the inverse of + * stdoutDemux.ts. + * + * PLOTAPP:\n a message + * PLOTBIN::\n
a binary frame + * + * The layout ipc.emit and anyplotlib's _binary_frame.encode_frame write on the + * Python side. hlen and plen are byte counts. The JSON text is JSON.stringify's + * own (compact, raw UTF-8); every decoder of the format accepts it. Imports + * nothing. + */ + +/** JSON.stringify replacer: a non-finite number is not JSON, so refuse it rather than emit null or NaN. */ +function refuseNonFinite(key: string, value: unknown): unknown { + if (typeof value === 'number' && !Number.isFinite(value)) { + throw new RangeError(`cannot frame the non-finite number ${value}${key ? ` at key "${key}"` : ''}`) + } + return value +} + +/** `PLOTAPP:\n`. JSON escapes every newline inside a string, so the message is one line. */ +export function encodeMessage(obj: Record): Buffer { + return Buffer.from(`PLOTAPP:${JSON.stringify(obj, refuseNonFinite)}\n`, 'utf8') +} + +/** + * The prefix line and the header bytes of a PLOTBIN frame. The payload is the + * caller's Buffer, written as its own chunk after these two, never copied. + */ +export function encodeBinary(header: Record, payload: Buffer): [Buffer, Buffer] { + const head = Buffer.from(JSON.stringify(header, refuseNonFinite), 'utf8') + return [Buffer.from(`PLOTBIN:${head.length}:${payload.length}\n`, 'ascii'), head] +} From bfbdd7b241b550a4ba3bd5e47b5d763ae056af59 Mon Sep 17 00:00:00 2001 From: Michael Spilman <130705596+TheDrOnos@users.noreply.github.com> Date: Fri, 25 Sep 2026 08:59:30 -0500 Subject: [PATCH 02/13] relay: TCP listener that binds exactly, splits client lines, and reports each close once --- de_shell/js/main/relay.test.ts | 388 +++++++++++++++++++++++++++++++++ de_shell/js/main/relay.ts | 153 +++++++++++++ 2 files changed, 541 insertions(+) create mode 100644 de_shell/js/main/relay.test.ts create mode 100644 de_shell/js/main/relay.ts diff --git a/de_shell/js/main/relay.test.ts b/de_shell/js/main/relay.test.ts new file mode 100644 index 0000000..e9b4471 --- /dev/null +++ b/de_shell/js/main/relay.test.ts @@ -0,0 +1,388 @@ +/** + * relay.test.ts — node:test suite for relay.ts on real loopback sockets + * (127.0.0.1, port 0), no Electron. The relay's timers run on an injected fake + * clock, so no test waits for a real timeout. + * + * Run: `node --test de_shell/js/main/relay.test.ts`, or via `npm run test:unit`. + */ +import { test } from 'node:test' +import assert from 'node:assert/strict' +import * as net from 'node:net' +import { createRelay } from './relay.ts' +import type { Relay, RelayCloseReason, RelayConnection, RelayOptions } from './relay.ts' + +type Ev = + | { kind: 'connection'; conn: number } + | { kind: 'line'; conn: number; line: string } + | { kind: 'close'; conn: number; reason: RelayCloseReason; err?: Error } + +interface Hooks { + onConnection?: (c: RelayConnection) => void + onLine?: (c: RelayConnection, line: string) => void +} + +interface Harness { + relay: Relay + events: Ev[] + conns: RelayConnection[] +} + +/** A relay on 127.0.0.1:0 that records every callback; `hooks` run after the record. */ +async function startRelay(extra: Partial = {}, hooks: Hooks = {}): Promise { + const events: Ev[] = [] + const conns: RelayConnection[] = [] + const relay = await createRelay({ + host: '127.0.0.1', + port: 0, + onConnection: (c) => { + conns.push(c) + events.push({ kind: 'connection', conn: conns.indexOf(c) }) + hooks.onConnection?.(c) + }, + onLine: (c, line) => { + events.push({ kind: 'line', conn: conns.indexOf(c), line }) + hooks.onLine?.(c, line) + }, + onClose: (c, reason, err) => { + events.push({ kind: 'close', conn: conns.indexOf(c), reason, err }) + }, + ...extra, + }) + return { relay, events, conns } +} + +function linesOf(events: Ev[], conn = 0): string[] { + return events.flatMap((e) => (e.kind === 'line' && e.conn === conn ? [e.line] : [])) +} + +function reasonsOf(events: Ev[], conn = 0): RelayCloseReason[] { + return events.flatMap((e) => (e.kind === 'close' && e.conn === conn ? [e.reason] : [])) +} + +function closeOf(events: Ev[], conn = 0): Extract | undefined { + return events.find( + (e): e is Extract => e.kind === 'close' && e.conn === conn, + ) +} + +function connectClient(port: number): Promise { + return new Promise((resolve, reject) => { + const s = net.connect({ host: '127.0.0.1', port }) + s.on('error', () => { /* resets on relay-initiated closes are expected */ }) + s.once('error', reject) + s.once('connect', () => resolve(s)) + }) +} + +function socketClosed(s: net.Socket): Promise { + return new Promise((resolve) => { + if (s.closed) resolve() + else s.once('close', () => resolve()) + }) +} + +const sleep = (ms: number): Promise => new Promise((r) => setTimeout(r, ms)) + +async function waitFor(pred: () => boolean, what: string, ms = 5000): Promise { + const t0 = Date.now() + while (!pred()) { + if (Date.now() - t0 > ms) throw new Error(`timed out waiting for ${what}`) + await sleep(10) + } +} + +interface Spied { + socket: net.Socket + log: string[] +} + +/** + * A createServer factory that records the server it made and every socket it + * accepted, with a log of the calls the relay makes on each socket + * ('setKeepAlive:true,15000', 'setNoDelay:true', 'cork', 'write:', …). + */ +function spyServer() { + const servers: net.Server[] = [] + const accepted: Spied[] = [] + const createServer = ((listener?: (s: net.Socket) => void): net.Server => { + const server = net.createServer(listener) + servers.push(server) + // Prepended so it runs before the relay's own connection listener. + server.prependListener('connection', (socket: net.Socket) => { + const log: string[] = [] + const target = socket as unknown as Record unknown> + for (const name of ['setKeepAlive', 'setNoDelay', 'cork', 'uncork', 'write', 'end', 'destroy']) { + const orig = target[name].bind(socket) + target[name] = (...a: unknown[]) => { + if (name === 'write') log.push(`write:${(a[0] as Buffer).length}`) + else log.push(a.length ? `${name}:${a.map(String).join(',')}` : name) + return orig(...a) + } + } + accepted.push({ socket, log }) + }) + return server + }) as typeof net.createServer + return { createServer, servers, accepted } +} + +test('reports the bound address and the ephemeral port', async () => { + const h = await startRelay() + try { + assert.equal(h.relay.address, '127.0.0.1') + assert.ok(h.relay.port > 0 && h.relay.port < 65536, `port ${h.relay.port}`) + const c = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 1, 'the connection') + c.destroy() + } finally { + await h.relay.close() + } +}) + +test('a bind to an address this machine lacks rejects with the OS error and leaves nothing listening', async () => { + // 192.0.2.1 is TEST-NET-1 (RFC 5737): assigned to no interface anywhere. + const spy = spyServer() + await assert.rejects( + createRelay({ + host: '192.0.2.1', + port: 0, + onConnection: () => {}, + onLine: () => {}, + onClose: () => {}, + createServer: spy.createServer, + }), + (err: unknown) => { + assert.equal((err as NodeJS.ErrnoException).code, 'EADDRNOTAVAIL') + return true + }, + ) + assert.equal(spy.servers.length, 1, 'the injected factory made the server') + assert.equal(spy.servers[0].listening, false) +}) + +test('a port already in use rejects with EADDRINUSE', async () => { + const first = await startRelay() + try { + await assert.rejects( + createRelay({ + host: '127.0.0.1', + port: first.relay.port, + onConnection: () => {}, + onLine: () => {}, + onClose: () => {}, + }), + (err: unknown) => { + assert.equal((err as NodeJS.ErrnoException).code, 'EADDRINUSE') + return true + }, + ) + } finally { + await first.relay.close() + } +}) + +test('the injected createServer makes the listening server', async () => { + const spy = spyServer() + const h = await startRelay({ createServer: spy.createServer }) + try { + assert.equal(spy.servers.length, 1) + assert.equal(spy.servers[0].listening, true) + assert.equal((spy.servers[0].address() as net.AddressInfo).port, h.relay.port) + const c = await connectClient(h.relay.port) + await waitFor(() => spy.accepted.length === 1 && h.conns.length === 1, 'the connection') + c.destroy() + } finally { + await h.relay.close() + } +}) + +test('every client gets onConnection with its remote address; the relay imposes no count', async () => { + const h = await startRelay() + const clients: net.Socket[] = [] + try { + for (let i = 0; i < 3; i++) clients.push(await connectClient(h.relay.port)) + await waitFor(() => h.conns.length === 3, 'three connections') + for (const c of h.conns) assert.equal(c.remoteAddress, '127.0.0.1') + assert.deepEqual( + new Set(h.conns.map((c) => c.remotePort)), + new Set(clients.map((c) => c.localPort)), + ) + assert.deepEqual(reasonsOf(h.events, 0), []) + } finally { + for (const c of clients) c.destroy() + await h.relay.close() + } +}) + +test('lines arrive whole across chunk splits, the first line included', async () => { + const h = await startRelay() + try { + const c = await connectClient(h.relay.port) + c.setNoDelay(true) + await waitFor(() => h.conns.length === 1, 'the connection') + // One byte per write, so the multi-byte ε is split too. + const bytes = Buffer.from('{"type":"hello","who":"εxx"}\n{"type":"action","name":"snap"}\n', 'utf8') + for (let i = 0; i < bytes.length; i++) { + c.write(bytes.subarray(i, i + 1)) + await sleep(1) + } + await waitFor(() => linesOf(h.events).length === 2, 'two lines') + assert.deepEqual(linesOf(h.events), [ + '{"type":"hello","who":"εxx"}', + '{"type":"action","name":"snap"}', + ]) + c.destroy() + } finally { + await h.relay.close() + } +}) + +test('\\r\\n is accepted and blank or whitespace-only lines are skipped', async () => { + const h = await startRelay() + try { + const c = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 1, 'the connection') + c.write('a\r\n\r\n \n\nb\n') + await waitFor(() => linesOf(h.events).length === 2, 'two lines') + await sleep(20) + assert.deepEqual(linesOf(h.events), ['a', 'b']) + c.destroy() + } finally { + await h.relay.close() + } +}) + +test('bytes that are not UTF-8 arrive as U+FFFD and the connection stays open', async () => { + const h = await startRelay() + try { + const c = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 1, 'the connection') + c.write(Buffer.from([0x66, 0xff, 0xfe, 0x0a])) + c.write('next\n') + await waitFor(() => linesOf(h.events).length === 2, 'two lines') + assert.deepEqual(linesOf(h.events), ['f\uFFFD\uFFFD', 'next']) + assert.deepEqual(reasonsOf(h.events), []) + c.destroy() + } finally { + await h.relay.close() + } +}) + +test('two clients are both delivered, each on its own connection', async () => { + const h = await startRelay() + try { + const a = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 1, 'first') + const b = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 2, 'second') + a.write('from a\n') + b.write('from b\n') + await waitFor( + () => linesOf(h.events, 0).length === 1 && linesOf(h.events, 1).length === 1, + 'both lines', + ) + assert.deepEqual(linesOf(h.events, 0), ['from a']) + assert.deepEqual(linesOf(h.events, 1), ['from b']) + a.destroy() + b.destroy() + } finally { + await h.relay.close() + } +}) + +test('a peer that half-closes: complete lines first, the unterminated tail dropped, then peer once', async () => { + const h = await startRelay() + try { + const c = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 1, 'the connection') + c.end('one\ntwo\npartial') + await waitFor(() => reasonsOf(h.events).length > 0, 'onClose') + await sleep(50) + assert.deepEqual(linesOf(h.events), ['one', 'two']) + assert.deepEqual(reasonsOf(h.events), ['peer']) + assert.equal(h.events.at(-1)?.kind, 'close', 'onClose comes after every line') + await socketClosed(c) + } finally { + await h.relay.close() + } +}) + +test('a socket error closes that connection as error with the error attached', async () => { + const spy = spyServer() + const h = await startRelay({ createServer: spy.createServer }) + try { + const c = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 1 && spy.accepted.length === 1, 'the connection') + spy.accepted[0].socket.destroy(new Error('injected')) + await waitFor(() => reasonsOf(h.events).length > 0, 'onClose') + await sleep(50) + assert.deepEqual(reasonsOf(h.events), ['error']) + assert.equal(closeOf(h.events)?.err?.message, 'injected') + await socketClosed(c) + } finally { + await h.relay.close() + } +}) + +test('an onLine that throws closes that connection as error; the relay keeps serving', async () => { + const boom = new Error('handler bug') + const h = await startRelay({}, { + onLine: (_c, line) => { + if (line === 'bad') throw boom + }, + }) + try { + const a = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 1, 'first') + a.write('bad\nnever delivered\n') + await waitFor(() => reasonsOf(h.events, 0).length === 1, 'onClose') + assert.deepEqual(reasonsOf(h.events, 0), ['error']) + assert.equal(closeOf(h.events, 0)?.err, boom) + assert.deepEqual(linesOf(h.events, 0), ['bad']) + const b = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 2, 'second') + b.write('fine\n') + await waitFor(() => linesOf(h.events, 1).length === 1, 'the second client line') + b.destroy() + } finally { + await h.relay.close() + } +}) + +test('an onConnection that throws closes that connection as error', async () => { + const h = await startRelay({}, { + onConnection: () => { + throw new Error('refused in handler') + }, + }) + try { + const c = await connectClient(h.relay.port) + await waitFor(() => reasonsOf(h.events).length === 1, 'onClose') + assert.deepEqual(reasonsOf(h.events), ['error']) + assert.equal(closeOf(h.events)?.err?.message, 'refused in handler') + await socketClosed(c) + } finally { + await h.relay.close() + } +}) + +test('close() ends every connection as app, stops the server, and resolves', async () => { + const spy = spyServer() + const h = await startRelay({ createServer: spy.createServer }) + const a = await connectClient(h.relay.port) + const b = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 2, 'two connections') + await h.relay.close() + assert.deepEqual(reasonsOf(h.events, 0), ['app']) + assert.deepEqual(reasonsOf(h.events, 1), ['app']) + assert.equal(spy.servers[0].listening, false) + await Promise.all([socketClosed(a), socketClosed(b)]) +}) + +test('close() with no connections resolves, and calling it again resolves too', async () => { + const spy = spyServer() + const { relay } = await startRelay({ createServer: spy.createServer }) + await Promise.all([relay.close(), relay.close()]) + await relay.close() + assert.equal(spy.servers[0].listening, false) +}) diff --git a/de_shell/js/main/relay.ts b/de_shell/js/main/relay.ts new file mode 100644 index 0000000..16b09c1 --- /dev/null +++ b/de_shell/js/main/relay.ts @@ -0,0 +1,153 @@ +/** + * relay.ts — a TCP listener that speaks the backend's framing to remote + * clients: an app's remote-control endpoint, a GUI on another machine. + * + * It moves bytes and nothing else. In: '\n'-terminated UTF-8 lines, handed to + * the app one string at a time. Who is admitted, what a line means and when a + * client is dropped are the app's decisions. No electron import, so it loads + * under `node --test` like backendProcess.ts. + */ +import * as net from 'node:net' + +export interface RelayConnection { + readonly remoteAddress: string + readonly remotePort: number +} + +export type RelayCloseReason = 'peer' | 'app' | 'hello-timeout' | 'line-too-long' | 'error' + +export interface Relay { + readonly address: string + readonly port: number + /** Ends every connection (reason 'app') and the server; resolves once the server has stopped. Safe to call twice. */ + close(): Promise +} + +export interface RelayOptions { + /** Bound exactly: no fallback, no discovery. '0.0.0.0' binds every interface only when the app passes it. */ + host: string + /** 0 = ephemeral; the bound port is reported on the Relay. */ + port: number + onConnection: (c: RelayConnection) => void + /** Every non-blank line, '\r\n' accepted, the first (hello) included. The relay never parses it. */ + onLine: (c: RelayConnection, line: string) => void + /** Exactly once per connection; `err` is set for 'error'. */ + onClose: (c: RelayConnection, reason: RelayCloseReason, err?: Error) => void + /** + * The TLS hook. Called once, as createServer(connectionListener), so a + * factory returning tls.createServer(tlsOptions, connectionListener) fits. + * Default net.createServer. + */ + createServer?: typeof net.createServer +} + +const NL = 0x0a + +function asError(e: unknown): Error { + return e instanceof Error ? e : new Error(String(e)) +} + +export function createRelay(opts: RelayOptions): Promise { + const makeServer = opts.createServer ?? net.createServer + const live = new Set<() => void>() // destroy-now closers of connections still open to the app + const sockets = new Set() // every accepted socket not yet 'close'd + + const accept = (socket: net.Socket): void => { + sockets.add(socket) + let closed = false + let pending: Buffer[] = [] // the partial line so far, as received + let pendingLen = 0 + + const conn: RelayConnection = { + remoteAddress: socket.remoteAddress ?? '', + remotePort: socket.remotePort ?? 0, + } + + const finish = (reason: RelayCloseReason, err?: Error): void => { + if (closed) return + closed = true + live.delete(destroyNow) + pending = [] + pendingLen = 0 + socket.destroy() + try { + opts.onClose(conn, reason, err) + } catch { + // Nowhere left to report it: the connection is already gone. + } + } + const destroyNow = (): void => finish('app') + + const onData = (chunk: Buffer): void => { + let start = 0 + while (!closed) { + const nl = chunk.indexOf(NL, start) + if (nl < 0) { + if (start < chunk.length) { + pending.push(chunk.subarray(start)) + pendingLen += chunk.length - start + } + return + } + let bytes = chunk.subarray(start, nl) + start = nl + 1 + if (pendingLen > 0) { + pending.push(bytes) + bytes = Buffer.concat(pending, pendingLen + bytes.length) + pending = [] + pendingLen = 0 + } + let line = bytes.toString('utf8') + if (line.endsWith('\r')) line = line.slice(0, -1) + if (line.trim() === '') continue + try { + opts.onLine(conn, line) + } catch (e) { + finish('error', asError(e)) + return + } + } + } + + socket.on('data', onData) + socket.on('end', () => finish('peer')) + socket.on('error', (err: Error) => finish('error', err)) + socket.on('close', () => { + sockets.delete(socket) + finish('peer') + }) + live.add(destroyNow) + try { + opts.onConnection(conn) + } catch (e) { + finish('error', asError(e)) + } + } + + const server = makeServer(accept) + return new Promise((resolve, reject) => { + const onListenError = (err: Error): void => reject(err) + server.once('error', onListenError) + server.listen(opts.port, opts.host, () => { + server.off('error', onListenError) + // After listen, a server 'error' is an accept failure (EMFILE and the + // like): transient, and the listener stays up. Handled so it cannot + // surface as an unhandled event in the main process. + server.on('error', () => {}) + const bound = server.address() as net.AddressInfo + let closing: Promise | null = null + resolve({ + address: bound.address, + port: bound.port, + close(): Promise { + closing ??= new Promise((done) => { + for (const destroyNow of [...live]) destroyNow() + for (const s of [...sockets]) s.destroy() + server.close(() => done()) + }) + return closing + }, + }) + }) + }) +} From 2db91e39934a0bd4558d31d253dbaa657e1bf509 Mon Sep 17 00:00:00 2001 From: Michael Spilman <130705596+TheDrOnos@users.noreply.github.com> Date: Fri, 25 Sep 2026 09:02:32 -0500 Subject: [PATCH 03/13] relay: pre-admission guard (hello timer, line caps, admit) on an injectable clock --- de_shell/js/main/relay.test.ts | 165 +++++++++++++++++++++++++++++++++ de_shell/js/main/relay.ts | 47 ++++++++++ 2 files changed, 212 insertions(+) diff --git a/de_shell/js/main/relay.test.ts b/de_shell/js/main/relay.test.ts index e9b4471..67873a6 100644 --- a/de_shell/js/main/relay.test.ts +++ b/de_shell/js/main/relay.test.ts @@ -386,3 +386,168 @@ test('close() with no connections resolves, and calling it again resolves too', await relay.close() assert.equal(spy.servers[0].listening, false) }) + +/** A manual clock for the relay's timers: nothing fires until advance(). */ +function fakeClock() { + let now = 0 + let nextId = 1 + const timers = new Map void }>() + return { + setTimeout: (fn: () => void, ms: number): unknown => { + const id = nextId++ + timers.set(id, { at: now + ms, fn }) + return id + }, + clearTimeout: (handle: unknown): void => { + timers.delete(handle as number) + }, + advance(ms: number): void { + now += ms + const due = [...timers].filter(([, t]) => t.at <= now).sort((x, y) => x[1].at - y[1].at) + for (const [id, t] of due) { + if (timers.delete(id)) t.fn() + } + }, + get pending(): number { + return timers.size + }, + } +} + +function clocked(clock: ReturnType, extra: Partial = {}): Partial { + return { setTimeout: clock.setTimeout, clearTimeout: clock.clearTimeout, ...extra } +} + +const admitOnHello: Hooks = { + onLine: (c, line) => { + if (line === 'hello') c.admit() + }, +} + +test('a silent client is closed as hello-timeout when the injected clock reaches helloTimeoutMs', async () => { + const clock = fakeClock() + const h = await startRelay(clocked(clock, { helloTimeoutMs: 10_000 })) + try { + const c = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 1, 'the connection') + clock.advance(9_999) + assert.deepEqual(reasonsOf(h.events), []) + clock.advance(1) + assert.deepEqual(reasonsOf(h.events), ['hello-timeout']) + await socketClosed(c) + } finally { + await h.relay.close() + } +}) + +test('the hello timer runs from accept to admit(): a line the app does not admit does not stop it', async () => { + const clock = fakeClock() + const h = await startRelay(clocked(clock, { helloTimeoutMs: 10_000 })) + try { + const c = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 1, 'the connection') + c.write('hello\n') + await waitFor(() => linesOf(h.events).length === 1, 'hello') + clock.advance(10_000) + assert.deepEqual(reasonsOf(h.events), ['hello-timeout']) + await socketClosed(c) + } finally { + await h.relay.close() + } +}) + +test('admit() stops the hello timer', async () => { + const clock = fakeClock() + const h = await startRelay(clocked(clock), admitOnHello) + try { + const c = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 1, 'the connection') + c.write('hello\n') + await waitFor(() => linesOf(h.events).length === 1, 'hello') + assert.equal(clock.pending, 0, 'admit() cancelled the timer') + clock.advance(1_000_000) + c.write('still here\n') + await waitFor(() => linesOf(h.events).length === 2, 'the second line') + assert.deepEqual(reasonsOf(h.events), []) + c.destroy() + } finally { + await h.relay.close() + } +}) + +test('a pre-admit line over the cap closes as line-too-long; the same line after admit() passes', async () => { + const clock = fakeClock() + const long = 'x'.repeat(17) + const h = await startRelay(clocked(clock, { preAdmitLineBytes: 16 }), admitOnHello) + try { + const a = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 1, 'first') + a.write(`${long}\n`) + await waitFor(() => reasonsOf(h.events, 0).length === 1, 'first onClose') + assert.deepEqual(reasonsOf(h.events, 0), ['line-too-long']) + assert.deepEqual(linesOf(h.events, 0), []) + const b = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 2, 'second') + // One write: admit() inside onLine('hello') lifts the cap for the rest of the same chunk. + b.write(`hello\n${long}\n`) + await waitFor(() => linesOf(h.events, 1).length === 2, 'both lines') + assert.deepEqual(linesOf(h.events, 1), ['hello', long]) + assert.deepEqual(reasonsOf(h.events, 1), []) + b.destroy() + } finally { + await h.relay.close() + } +}) + +test('a line of exactly the cap passes; a pre-admit partial line past it closes as line-too-long', async () => { + const clock = fakeClock() + const h = await startRelay(clocked(clock, { preAdmitLineBytes: 16 })) + try { + const c = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 1, 'the connection') + c.write(`${'y'.repeat(16)}\n`) + await waitFor(() => linesOf(h.events).length === 1, 'the 16-byte line') + c.write('z'.repeat(10)) + await sleep(50) + assert.deepEqual(reasonsOf(h.events), [], 'ten bytes are under the cap') + c.write('z'.repeat(7)) // no newline ever: the partial line alone crosses the cap + await waitFor(() => reasonsOf(h.events).length === 1, 'onClose') + assert.deepEqual(reasonsOf(h.events), ['line-too-long']) + assert.deepEqual(linesOf(h.events), ['y'.repeat(16)]) + await socketClosed(c) + } finally { + await h.relay.close() + } +}) + +test('after admit() the cap is lineBytes', async () => { + const clock = fakeClock() + const h = await startRelay(clocked(clock, { preAdmitLineBytes: 16, lineBytes: 32 }), admitOnHello) + try { + const c = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 1, 'the connection') + c.write(`hello\n${'y'.repeat(32)}\n${'y'.repeat(33)}\n`) + await waitFor(() => reasonsOf(h.events).length === 1, 'onClose') + assert.deepEqual(linesOf(h.events), ['hello', 'y'.repeat(32)]) + assert.deepEqual(reasonsOf(h.events), ['line-too-long']) + } finally { + await h.relay.close() + } +}) + +test('a connection that closes cancels its hello timer', async () => { + const clock = fakeClock() + const h = await startRelay(clocked(clock)) + try { + const c = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 1, 'the connection') + assert.equal(clock.pending, 1) + c.end() + await waitFor(() => reasonsOf(h.events).length === 1, 'onClose') + assert.equal(clock.pending, 0) + clock.advance(1_000_000) + assert.deepEqual(reasonsOf(h.events), ['peer']) + } finally { + await h.relay.close() + } +}) diff --git a/de_shell/js/main/relay.ts b/de_shell/js/main/relay.ts index 16b09c1..f9763e7 100644 --- a/de_shell/js/main/relay.ts +++ b/de_shell/js/main/relay.ts @@ -6,12 +6,19 @@ * the app one string at a time. Who is admitted, what a line means and when a * client is dropped are the app's decisions. No electron import, so it loads * under `node --test` like backendProcess.ts. + * + * Until the app calls admit(), a connection is on probation: it is closed as + * 'hello-timeout' helloTimeoutMs after accept, and as 'line-too-long' once a + * line (or a partial one) passes preAdmitLineBytes. After admit() the cap is + * lineBytes and there is no timer. */ import * as net from 'node:net' export interface RelayConnection { readonly remoteAddress: string readonly remotePort: number + /** Stops the hello timer and lifts the line cap from preAdmitLineBytes to lineBytes. */ + admit(): void } export type RelayCloseReason = 'peer' | 'app' | 'hello-timeout' | 'line-too-long' | 'error' @@ -33,12 +40,22 @@ export interface RelayOptions { onLine: (c: RelayConnection, line: string) => void /** Exactly once per connection; `err` is set for 'error'. */ onClose: (c: RelayConnection, reason: RelayCloseReason, err?: Error) => void + /** Accept to admit(); past it the connection is closed as 'hello-timeout'. Default 10 000. */ + helloTimeoutMs?: number + /** Line cap in bytes (those before the '\n') until admit(). Default 64 KiB. */ + preAdmitLineBytes?: number + /** Line cap in bytes after admit(). Default 16 MiB. */ + lineBytes?: number /** * The TLS hook. Called once, as createServer(connectionListener), so a * factory returning tls.createServer(tlsOptions, connectionListener) fits. * Default net.createServer. */ createServer?: typeof net.createServer + /** The clock for the relay's timers; tests inject a fake one. Default: the global setTimeout. */ + setTimeout?: (fn: () => void, ms: number) => unknown + /** Pairs with setTimeout. Default: the global clearTimeout. */ + clearTimeout?: (handle: unknown) => void } const NL = 0x0a @@ -48,24 +65,46 @@ function asError(e: unknown): Error { } export function createRelay(opts: RelayOptions): Promise { + const helloTimeoutMs = opts.helloTimeoutMs ?? 10_000 + const preAdmitLineBytes = opts.preAdmitLineBytes ?? 64 * 1024 + const lineBytes = opts.lineBytes ?? 16 * 1024 * 1024 const makeServer = opts.createServer ?? net.createServer + const arm = opts.setTimeout ?? ((fn: () => void, ms: number): unknown => setTimeout(fn, ms)) + const disarm = opts.clearTimeout + ?? ((handle: unknown): void => clearTimeout(handle as Parameters[0])) const live = new Set<() => void>() // destroy-now closers of connections still open to the app const sockets = new Set() // every accepted socket not yet 'close'd const accept = (socket: net.Socket): void => { sockets.add(socket) + let admitted = false let closed = false + let hello: unknown = null let pending: Buffer[] = [] // the partial line so far, as received let pendingLen = 0 + const cap = (): number => (admitted ? lineBytes : preAdmitLineBytes) + const stopHello = (): void => { + if (hello !== null) { + disarm(hello) + hello = null + } + } + const conn: RelayConnection = { remoteAddress: socket.remoteAddress ?? '', remotePort: socket.remotePort ?? 0, + admit(): void { + if (closed || admitted) return + admitted = true + stopHello() + }, } const finish = (reason: RelayCloseReason, err?: Error): void => { if (closed) return closed = true + stopHello() live.delete(destroyNow) pending = [] pendingLen = 0 @@ -87,6 +126,8 @@ export function createRelay(opts: RelayOptions): Promise { pending.push(chunk.subarray(start)) pendingLen += chunk.length - start } + // A client that never sends '\n' is capped here, not buffered without bound. + if (pendingLen > cap()) finish('line-too-long') return } let bytes = chunk.subarray(start, nl) @@ -97,6 +138,11 @@ export function createRelay(opts: RelayOptions): Promise { pending = [] pendingLen = 0 } + // Checked per line, so an admit() inside onLine lifts the cap for the rest of this chunk. + if (bytes.length > cap()) { + finish('line-too-long') + return + } let line = bytes.toString('utf8') if (line.endsWith('\r')) line = line.slice(0, -1) if (line.trim() === '') continue @@ -117,6 +163,7 @@ export function createRelay(opts: RelayOptions): Promise { finish('peer') }) live.add(destroyNow) + hello = arm(() => finish('hello-timeout'), helloTimeoutMs) try { opts.onConnection(conn) } catch (e) { From 5b91002a24bdcea9fce0356560ffd082f29a3c49 Mon Sep 17 00:00:00 2001 From: Michael Spilman <130705596+TheDrOnos@users.noreply.github.com> Date: Fri, 25 Sep 2026 09:05:15 -0500 Subject: [PATCH 04/13] relay: framed writes (corked binary), graceful close, keepalive and no-delay --- de_shell/js/main/relay.test.ts | 244 ++++++++++++++++++++++++++++++++- de_shell/js/main/relay.ts | 65 ++++++++- 2 files changed, 304 insertions(+), 5 deletions(-) diff --git a/de_shell/js/main/relay.test.ts b/de_shell/js/main/relay.test.ts index 67873a6..24a0e4c 100644 --- a/de_shell/js/main/relay.test.ts +++ b/de_shell/js/main/relay.test.ts @@ -8,7 +8,8 @@ import { test } from 'node:test' import assert from 'node:assert/strict' import * as net from 'node:net' -import { createRelay } from './relay.ts' +import { CLOSE_GRACE_MS, createRelay } from './relay.ts' +import { createStdoutDemux } from './stdoutDemux.ts' import type { Relay, RelayCloseReason, RelayConnection, RelayOptions } from './relay.ts' type Ev = @@ -551,3 +552,244 @@ test('a connection that closes cancels its hello timer', async () => { await h.relay.close() } }) + +type Unit = + | { kind: 'message'; msg: Record } + | { kind: 'binary'; header: Record; payload: Buffer } + | { kind: 'stream'; text: string } + +/** Decode everything the relay sends to `c` with the backend's own demuxer. */ +function collect(c: net.Socket): Unit[] { + const units: Unit[] = [] + const demux = createStdoutDemux({ + onMessage: (msg) => units.push({ kind: 'message', msg }), + onStream: (text) => units.push({ kind: 'stream', text }), + onBinary: (header, payload) => units.push({ kind: 'binary', header, payload }), + }) + c.on('data', (chunk: Buffer) => demux.push(chunk)) + return units +} + +/** Write 1 MiB frames until the socket queues bytes it cannot hand to the kernel (the peer is not reading). */ +function fillUntilQueued(conn: RelayConnection): void { + const block = Buffer.alloc(1 << 20) + for (let i = 0; i < 256 && conn.writableLength === 0; i++) { + conn.writeBinary({ fig_id: 'fill', key: 'k', i }, block) + } + assert.ok(conn.writableLength > 0, 'the socket never queued; is the peer reading?') +} + +test('writeMessage and writeBinary arrive in order and decode with the demuxer', async () => { + const nasty = Buffer.concat([Buffer.from('\nPLOTBIN:9:9\nPLOTAPP:{}\n', 'ascii'), Buffer.alloc(4096, 7)]) + const h = await startRelay({}, { + onConnection: (conn) => { + conn.writeMessage({ type: 'welcome', v: 1 }) + conn.writeBinary({ fig_id: 'f1', key: 'image', label: 'εxx' }, nasty) + conn.writeMessage({ type: 'done' }) + }, + }) + try { + const c = await connectClient(h.relay.port) + const units = collect(c) + await waitFor(() => units.length === 3, 'three units') + assert.deepEqual(units, [ + { kind: 'message', msg: { type: 'welcome', v: 1 } }, + { kind: 'binary', header: { fig_id: 'f1', key: 'image', label: 'εxx' }, payload: nasty }, + { kind: 'message', msg: { type: 'done' } }, + ]) + c.destroy() + } finally { + await h.relay.close() + } +}) + +test('writeBinary corks the socket around its three writes', async () => { + const spy = spyServer() + const h = await startRelay({ createServer: spy.createServer }) + try { + const c = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 1 && spy.accepted.length === 1, 'the connection') + const log = spy.accepted[0].log + const from = log.length + const header = { fig_id: 'f', key: 'k' } + h.conns[0].writeBinary(header, Buffer.alloc(1000, 1)) + const hlen = Buffer.byteLength(JSON.stringify(header)) + assert.deepEqual(log.slice(from), [ + 'cork', + `write:${`PLOTBIN:${hlen}:1000\n`.length}`, + `write:${hlen}`, + 'write:1000', + 'uncork', + ]) + c.destroy() + } finally { + await h.relay.close() + } +}) + +test('writeMessage refuses a non-finite number and writes nothing; the connection stays usable', async () => { + const h = await startRelay() + try { + const c = await connectClient(h.relay.port) + const units = collect(c) + await waitFor(() => h.conns.length === 1, 'the connection') + const conn = h.conns[0] + assert.throws(() => conn.writeMessage({ type: 'fit', value: NaN }), RangeError) + assert.throws( + () => conn.writeBinary({ fig_id: 'f', key: 'k', clim: [0, Infinity] }, Buffer.alloc(8)), + RangeError, + ) + conn.writeMessage({ type: 'after' }) + await waitFor(() => units.length === 1, 'one unit') + await sleep(20) + assert.deepEqual(units, [{ kind: 'message', msg: { type: 'after' } }]) + assert.deepEqual(reasonsOf(h.events), []) + c.destroy() + } finally { + await h.relay.close() + } +}) + +test('everything written before close() arrives before EOF, the refusal last; onClose reports app once', async () => { + // 8 MiB ahead of the refusal, so bytes are still queued in the relay when + // close() runs: a close that destroyed the socket would lose them (a lone + // small write reaches the kernel at once and would not tell the two apart). + const block = Buffer.alloc(1 << 20, 3) + let queuedAtClose = -1 + const h = await startRelay({}, { + onLine: (conn, line) => { + if (line === 'hello') { + for (let i = 0; i < 8; i++) conn.writeBinary({ fig_id: 'f', key: 'k', i }, block) + conn.writeMessage({ type: 'refused', reason: 'not paired' }) + queuedAtClose = conn.writableLength + conn.close() + } + }, + }) + try { + const c = await connectClient(h.relay.port) + const units = collect(c) + c.write('hello\n') + await socketClosed(c) + assert.ok(queuedAtClose > 0, 'bytes were still queued in the relay when close() was called') + assert.equal(units.length, 9) + assert.deepEqual( + units.slice(0, 8).map((u) => (u.kind === 'binary' ? [u.header.i, u.payload.length] : null)), + [0, 1, 2, 3, 4, 5, 6, 7].map((i) => [i, 1 << 20]), + ) + assert.deepEqual(units[8], { kind: 'message', msg: { type: 'refused', reason: 'not paired' } }) + await sleep(50) + assert.deepEqual(reasonsOf(h.events), ['app']) + } finally { + await h.relay.close() + } +}) + +test('close() to a client that never reads destroys the socket after CLOSE_GRACE_MS', async () => { + const clock = fakeClock() + const spy = spyServer() + const h = await startRelay(clocked(clock, { createServer: spy.createServer })) + const c = await connectClient(h.relay.port) + c.pause() + try { + await waitFor(() => h.conns.length === 1 && spy.accepted.length === 1, 'the connection') + const conn = h.conns[0] + conn.admit() + fillUntilQueued(conn) + conn.close() + assert.deepEqual(reasonsOf(h.events), ['app']) + const socket = spy.accepted[0].socket + clock.advance(CLOSE_GRACE_MS - 1) + assert.equal(socket.destroyed, false, 'still flushing inside the grace period') + clock.advance(1) + assert.equal(socket.destroyed, true) + assert.deepEqual(reasonsOf(h.events), ['app'], 'onClose fired once') + } finally { + await h.relay.close() + c.destroy() + } +}) + +test('writes after close() are dropped without throwing', async () => { + const h = await startRelay() + try { + const c = await connectClient(h.relay.port) + const units = collect(c) + await waitFor(() => h.conns.length === 1, 'the connection') + const conn = h.conns[0] + conn.close() + conn.writeMessage({ type: 'late' }) + conn.writeBinary({ fig_id: 'f', key: 'k' }, Buffer.alloc(4)) + conn.close() + conn.admit() + await socketClosed(c) + assert.deepEqual(units, []) + assert.deepEqual(reasonsOf(h.events), ['app']) + } finally { + await h.relay.close() + } +}) + +test('writableLength grows when the peer stops reading', async () => { + const h = await startRelay() + const c = await connectClient(h.relay.port) + c.pause() + try { + await waitFor(() => h.conns.length === 1, 'the connection') + assert.equal(h.conns[0].writableLength, 0) + fillUntilQueued(h.conns[0]) + } finally { + await h.relay.close() + c.destroy() + } +}) + +test('keepalive and no-delay are set on each accepted socket; keepAliveMs 0 leaves keepalive off', async () => { + const cases: Array<[number | undefined, string | null]> = [ + [undefined, 'setKeepAlive:true,15000'], + [2500, 'setKeepAlive:true,2500'], + [0, null], + ] + for (const [keepAliveMs, expected] of cases) { + const spy = spyServer() + const extra: Partial = { createServer: spy.createServer } + if (keepAliveMs !== undefined) extra.keepAliveMs = keepAliveMs + const h = await startRelay(extra) + try { + const c = await connectClient(h.relay.port) + await waitFor(() => spy.accepted.length === 1, 'the connection') + const log = spy.accepted[0].log + assert.ok(log.includes('setNoDelay:true'), log.join(' ')) + assert.deepEqual( + log.filter((l) => l.startsWith('setKeepAlive')), + expected === null ? [] : [expected], + `keepAliveMs ${String(keepAliveMs)}`, + ) + c.destroy() + } finally { + await h.relay.close() + } + } +}) + +test('relay.close() while a connection is still flushing its close: both resolve, onClose once each', async () => { + const clock = fakeClock() + const spy = spyServer() + const h = await startRelay(clocked(clock, { createServer: spy.createServer })) + const a = await connectClient(h.relay.port) + a.pause() + await waitFor(() => h.conns.length === 1, 'first') + const b = await connectClient(h.relay.port) + await waitFor(() => h.conns.length === 2, 'second') + h.conns[0].admit() + fillUntilQueued(h.conns[0]) + h.conns[0].close() // flushing to a peer that will not read + await Promise.all([h.relay.close(), h.relay.close()]) + assert.deepEqual(reasonsOf(h.events, 0), ['app']) + assert.deepEqual(reasonsOf(h.events, 1), ['app']) + assert.equal(spy.accepted[0].socket.destroyed, true) + assert.equal(spy.servers[0].listening, false) + await waitFor(() => clock.pending === 0, 'the grace and hello timers cancelled') + a.destroy() + b.destroy() +}) diff --git a/de_shell/js/main/relay.ts b/de_shell/js/main/relay.ts index f9763e7..ab414cd 100644 --- a/de_shell/js/main/relay.ts +++ b/de_shell/js/main/relay.ts @@ -3,7 +3,8 @@ * clients: an app's remote-control endpoint, a GUI on another machine. * * It moves bytes and nothing else. In: '\n'-terminated UTF-8 lines, handed to - * the app one string at a time. Who is admitted, what a line means and when a + * the app one string at a time. Out: PLOTAPP messages and PLOTBIN frames, + * built by framing.ts. Who is admitted, what a line means and when a slow * client is dropped are the app's decisions. No electron import, so it loads * under `node --test` like backendProcess.ts. * @@ -11,14 +12,32 @@ * 'hello-timeout' helloTimeoutMs after accept, and as 'line-too-long' once a * line (or a partial one) passes preAdmitLineBytes. After admit() the cap is * lineBytes and there is no timer. + * + * Closing: onClose fires exactly once per connection, when the relay or the + * app ends it. conn.close() flushes what the app already wrote and then sends + * FIN, so a refusal written just before it still arrives; a peer that will not + * read is destroyed after CLOSE_GRACE_MS. Every other close destroys at once. + * + * Keepalive: Windows honours the initial delay only; the probe interval is the + * OS default. The app's own supersede rule is the half-open recovery on every + * platform. */ import * as net from 'node:net' +import { encodeBinary, encodeMessage } from './framing.ts' export interface RelayConnection { readonly remoteAddress: string readonly remotePort: number + /** Bytes queued on the socket and not yet flushed. The relay never buffers on the app's behalf. */ + readonly writableLength: number /** Stops the hello timer and lifts the line cap from preAdmitLineBytes to lineBytes. */ admit(): void + /** Throws on a non-finite number, having written nothing. Dropped once the connection is closed. */ + writeMessage(obj: Record): void + /** Prefix, header and payload go out corked, together. The payload is not copied. Dropped once closed. */ + writeBinary(header: Record, payload: Buffer): void + /** Flush what was written, then FIN; onClose(..., 'app') fires now. */ + close(): void } export type RelayCloseReason = 'peer' | 'app' | 'hello-timeout' | 'line-too-long' | 'error' @@ -46,6 +65,8 @@ export interface RelayOptions { preAdmitLineBytes?: number /** Line cap in bytes after admit(). Default 16 MiB. */ lineBytes?: number + /** TCP keepalive initial delay on every accepted socket; 0 leaves keepalive off. Default 15 000. */ + keepAliveMs?: number /** * The TLS hook. Called once, as createServer(connectionListener), so a * factory returning tls.createServer(tlsOptions, connectionListener) fits. @@ -58,6 +79,9 @@ export interface RelayOptions { clearTimeout?: (handle: unknown) => void } +/** How long conn.close() waits for queued writes to flush before destroying the socket. */ +export const CLOSE_GRACE_MS = 5000 + const NL = 0x0a function asError(e: unknown): Error { @@ -68,18 +92,20 @@ export function createRelay(opts: RelayOptions): Promise { const helloTimeoutMs = opts.helloTimeoutMs ?? 10_000 const preAdmitLineBytes = opts.preAdmitLineBytes ?? 64 * 1024 const lineBytes = opts.lineBytes ?? 16 * 1024 * 1024 + const keepAliveMs = opts.keepAliveMs ?? 15_000 const makeServer = opts.createServer ?? net.createServer const arm = opts.setTimeout ?? ((fn: () => void, ms: number): unknown => setTimeout(fn, ms)) const disarm = opts.clearTimeout ?? ((handle: unknown): void => clearTimeout(handle as Parameters[0])) const live = new Set<() => void>() // destroy-now closers of connections still open to the app - const sockets = new Set() // every accepted socket not yet 'close'd + const sockets = new Set() // every accepted socket not yet 'close'd, graceful closes included const accept = (socket: net.Socket): void => { sockets.add(socket) let admitted = false let closed = false let hello: unknown = null + let grace: unknown = null let pending: Buffer[] = [] // the partial line so far, as received let pendingLen = 0 @@ -94,21 +120,46 @@ export function createRelay(opts: RelayOptions): Promise { const conn: RelayConnection = { remoteAddress: socket.remoteAddress ?? '', remotePort: socket.remotePort ?? 0, + get writableLength(): number { + return socket.writableLength + }, admit(): void { if (closed || admitted) return admitted = true stopHello() }, + writeMessage(obj: Record): void { + const bytes = encodeMessage(obj) + if (closed) return + socket.write(bytes) + }, + writeBinary(header: Record, payload: Buffer): void { + const [prefix, head] = encodeBinary(header, payload) + if (closed) return + socket.cork() + socket.write(prefix) + socket.write(head) + if (payload.length > 0) socket.write(payload) + socket.uncork() + }, + close(): void { + finish('app', undefined, true) + }, } - const finish = (reason: RelayCloseReason, err?: Error): void => { + const finish = (reason: RelayCloseReason, err?: Error, graceful = false): void => { if (closed) return closed = true stopHello() live.delete(destroyNow) pending = [] pendingLen = 0 - socket.destroy() + if (graceful) { + socket.end() + grace = arm(() => socket.destroy(), CLOSE_GRACE_MS) + } else { + socket.destroy() + } try { opts.onClose(conn, reason, err) } catch { @@ -160,8 +211,14 @@ export function createRelay(opts: RelayOptions): Promise { socket.on('error', (err: Error) => finish('error', err)) socket.on('close', () => { sockets.delete(socket) + if (grace !== null) { + disarm(grace) + grace = null + } finish('peer') }) + if (keepAliveMs > 0) socket.setKeepAlive(true, keepAliveMs) + socket.setNoDelay(true) // small request/reply lines must not wait on Nagle + delayed ACK live.add(destroyNow) hello = arm(() => finish('hello-timeout'), helloTimeoutMs) try { From 088e2dae4f7c50b4e2bfcfc2e5c40b9b49420b74 Mon Sep 17 00:00:00 2001 From: Michael Spilman <130705596+TheDrOnos@users.noreply.github.com> Date: Fri, 25 Sep 2026 09:07:59 -0500 Subject: [PATCH 05/13] shell-main: export createRelay and the framing encoders --- de_shell/js/main/index.ts | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/de_shell/js/main/index.ts b/de_shell/js/main/index.ts index d32b03c..8d59c30 100644 --- a/de_shell/js/main/index.ts +++ b/de_shell/js/main/index.ts @@ -22,6 +22,10 @@ export { export { recordProblem, recordedProblems } from './problemLog' export type { Problem } from './problemLog' +export { createRelay } from './relay' +export type { Relay, RelayConnection, RelayCloseReason, RelayOptions } from './relay' +export { encodeMessage, encodeBinary } from './framing' + export { initErrorReporting, reportingConfigured, collectDiagnostics, submitReport, } from './errorReport' From 8eb262f81934feccdf43f05a39d08918ae236b82 Mon Sep 17 00:00:00 2001 From: Michael Spilman <130705596+TheDrOnos@users.noreply.github.com> Date: Fri, 25 Sep 2026 09:09:49 -0500 Subject: [PATCH 06/13] remote_client: stdlib decoder for the PLOTAPP/PLOTBIN stream, pinned against the real writers --- de_shell/remote_client.py | 103 +++++++++++++++++++++++++ tests/test_remote_client.py | 149 ++++++++++++++++++++++++++++++++++++ 2 files changed, 252 insertions(+) create mode 100644 de_shell/remote_client.py create mode 100644 tests/test_remote_client.py diff --git a/de_shell/remote_client.py b/de_shell/remote_client.py new file mode 100644 index 0000000..d811746 --- /dev/null +++ b/de_shell/remote_client.py @@ -0,0 +1,103 @@ +""" +remote_client.py — the Python end of a de-shell relay (``de_shell/js/main/relay.ts``). + +The relay speaks the backend's own framing to a remote peer: the peer sends +bare JSON lines; the relay sends ``PLOTAPP:\\n`` messages and +``PLOTBIN::\\n
`` binary frames. This module +decodes that stream. + +Standard library only: importable without numpy, anyplotlib, asyncio or +Electron (tests/test_remote_client.py checks it in a clean interpreter). +""" +from __future__ import annotations + +import json +import re + +_PLOTAPP = "PLOTAPP:" +_PLOTBIN = b"PLOTBIN:" +_PREFIX = re.compile(rb"PLOTBIN:(\d+):(\d+)") + + +def _malformed(what: str, line: str | bytes) -> ValueError: + """A ValueError naming the unit, with the offending text or bytes on ``.line``.""" + err = ValueError(f"malformed {what}: {line[:200]!r}") + err.line = line + return err + + +class _Decoder: + """Bytes in, protocol units out. + + Mirrors ``stdoutDemux.ts``: the unit trace does not depend on how the bytes + are chunked, and a frame waits for all ``hlen + plen`` bytes. One + difference: a malformed unit raises ``ValueError`` (after it is consumed, + so the next ``pop`` continues) instead of being dropped. + """ + + def __init__(self) -> None: + self._buf = bytearray() + self._scan = 0 # self._buf[:self._scan] holds no b"\n": never rescan it + + @property + def buffered(self) -> int: + """Bytes received and not yet part of a returned unit.""" + return len(self._buf) + + def feed(self, data: bytes) -> None: + self._buf += data + + def pop(self) -> tuple | None: + """The next complete unit, or None until more bytes arrive. + + ``("message", dict)`` for a PLOTAPP line, ``("binary", header, payload)`` + for a PLOTBIN frame, ``("stream", str)`` for any other non-blank line + (without its line ending). Blank lines are skipped. + """ + buf = self._buf + while buf: + nl = buf.find(b"\n", self._scan) + if nl < 0: + self._scan = len(buf) + return None + # A buffer that starts with the 8-byte marker has no newline before + # byte 8, so a newline found earlier always ends a text line. + if buf.startswith(_PLOTBIN): + prefix = bytes(buf[:nl]) + match = _PREFIX.fullmatch(prefix) + if match is None: + self._consume(nl + 1) + raise _malformed("PLOTBIN prefix", prefix) + hlen, plen = int(match.group(1)), int(match.group(2)) + end = nl + 1 + hlen + plen + if len(buf) < end: + return None + raw_header = bytes(buf[nl + 1:nl + 1 + hlen]) + payload = bytes(buf[nl + 1 + hlen:end]) + self._consume(end) + try: + header = json.loads(raw_header.decode("utf-8")) + except ValueError: # UnicodeDecodeError and JSONDecodeError alike + raise _malformed("PLOTBIN header", raw_header) from None + if not isinstance(header, dict): + raise _malformed("PLOTBIN header", raw_header) + return ("binary", header, payload) + line = bytes(buf[:nl]).decode("utf-8", errors="replace") + self._consume(nl + 1) + if line.endswith("\r"): + line = line[:-1] + if line.startswith(_PLOTAPP): + try: + msg = json.loads(line[len(_PLOTAPP):]) + except ValueError: + raise _malformed("PLOTAPP message", line) from None + if not isinstance(msg, dict): + raise _malformed("PLOTAPP message", line) + return ("message", msg) + if line.strip(): + return ("stream", line) + return None + + def _consume(self, n: int) -> None: + del self._buf[:n] + self._scan = 0 diff --git a/tests/test_remote_client.py b/tests/test_remote_client.py new file mode 100644 index 0000000..4d2ce24 --- /dev/null +++ b/tests/test_remote_client.py @@ -0,0 +1,149 @@ +""" +test_remote_client.py — the synchronous relay client. + +The decoder is fed bytes written by the REAL Python writers (ipc.emit, +ipc._write_line, and anyplotlib's encode_frame through ipc._write_binary) +across a socketpair, whole and one byte at a time, and must produce the same +trace both ways. +""" +from __future__ import annotations + +import json +import socket +import subprocess +import sys + +import pytest +from anyplotlib._binary_frame import encode_frame + +from de_shell import ipc +from de_shell.remote_client import _Decoder + +NASTY = b"\nPLOTBIN:9:9\nPLOTAPP:{}\n" + bytes(i & 0xFF for i in range(3000)) + + +def _through_socketpair(monkeypatch, write) -> bytes: + """Point ipc's protocol channel at one end of a socketpair, run ``write()``, + and return every byte that came out of the other end. Kept a few KiB, well + under any socketpair buffer, so nothing needs a reader thread.""" + a, b = socket.socketpair() + try: + out = a.makefile("w", encoding="utf-8", newline="\n") + monkeypatch.setattr(ipc, "_PROTOCOL_OUT", out) + write() + out.flush() + out.close() + a.shutdown(socket.SHUT_WR) + b.settimeout(5.0) + data = bytearray() + while chunk := b.recv(65536): + data += chunk + return bytes(data) + finally: + a.close() + b.close() + + +def _write_session() -> None: + ipc._write_line("starting up\n") + ipc.emit({"type": "status", "text": "Cluster ready"}) + ipc._write_binary(encode_frame("f1", "image", {"dims": [2, 1500], "dtype": "uint8"}, NASTY)) + ipc.emit({"type": "fit", "label": "εxx Å", "value": float("nan")}) # ipc sanitizes NaN to null + ipc._write_line("\n") + ipc._write_line(" \n") + ipc._write_line("mid-run log\n") + ipc._write_binary(encode_frame("f2", "spec", {}, b"")) + ipc.emit({"type": "done"}) + + +EXPECTED = [ + ("stream", "starting up"), + ("message", {"type": "status", "text": "Cluster ready"}), + ("binary", {"dims": [2, 1500], "dtype": "uint8", "fig_id": "f1", "key": "image"}, NASTY), + ("message", {"type": "fit", "label": "εxx Å", "value": None}), + ("stream", "mid-run log"), + ("binary", {"fig_id": "f2", "key": "spec"}, b""), + ("message", {"type": "done"}), +] + + +def _trace(data: bytes, step: int) -> list: + decoder = _Decoder() + units = [] + for pos in range(0, len(data), step): + decoder.feed(data[pos:pos + step]) + while (unit := decoder.pop()) is not None: + units.append(unit) + assert decoder.buffered == 0, "bytes left over after the last unit" + return units + + +def test_the_writers_output_decodes_to_the_expected_trace(monkeypatch): + data = _through_socketpair(monkeypatch, _write_session) + assert _trace(data, len(data)) == EXPECTED + + +def test_the_trace_is_the_same_one_byte_at_a_time(monkeypatch): + data = _through_socketpair(monkeypatch, _write_session) + assert _trace(data, 1) == EXPECTED + assert _trace(data, 7) == EXPECTED + + +def test_a_frame_split_at_every_byte_reassembles(): + frame = encode_frame("x", "image", {"label": "εxx"}, NASTY[:64]) + whole = [("binary", {"label": "εxx", "fig_id": "x", "key": "image"}, NASTY[:64])] + for cut in range(1, len(frame)): + decoder = _Decoder() + units = [] + for part in (frame[:cut], frame[cut:]): + decoder.feed(part) + while (unit := decoder.pop()) is not None: + units.append(unit) + assert units == whole, f"split at byte {cut}" + + +def test_a_malformed_prefix_raises_with_the_line_and_the_next_unit_decodes(): + decoder = _Decoder() + decoder.feed(b'PLOTBIN:12:x\nPLOTAPP:{"type":"after"}\n') + with pytest.raises(ValueError, match="PLOTBIN prefix") as excinfo: + decoder.pop() + assert excinfo.value.line == b"PLOTBIN:12:x" + assert decoder.pop() == ("message", {"type": "after"}) + + +def test_bad_json_after_plotapp_raises_with_the_line_and_the_next_unit_decodes(): + decoder = _Decoder() + decoder.feed(b'PLOTAPP:{nope\nPLOTAPP:{"type":"after"}\n') + with pytest.raises(ValueError, match="PLOTAPP message") as excinfo: + decoder.pop() + assert excinfo.value.line == "PLOTAPP:{nope" + assert decoder.pop() == ("message", {"type": "after"}) + + +def test_a_malformed_header_raises_after_consuming_the_frame(): + decoder = _Decoder() + decoder.feed(b"PLOTBIN:7:3\n{brokenabc" + b'PLOTAPP:{"type":"after"}\n') + with pytest.raises(ValueError, match="PLOTBIN header") as excinfo: + decoder.pop() + assert excinfo.value.line == b"{broken" + assert decoder.pop() == ("message", {"type": "after"}) + + +def test_crlf_is_stripped_and_bytes_that_are_not_utf8_become_replacement_characters(): + decoder = _Decoder() + decoder.feed(b"log line\r\n" + b"f\xff\xfe\n") + assert decoder.pop() == ("stream", "log line") + assert decoder.pop() == ("stream", "f\ufffd\ufffd") + assert decoder.pop() is None + + +def test_the_client_imports_only_the_standard_library(): + probe = ( + "import json, sys\n" + "import de_shell.remote_client\n" + "heavy = ('numpy', 'anyplotlib', 'asyncio', 'yaml')\n" + "print(json.dumps(sorted(m for m in heavy if m in sys.modules)))\n" + ) + proc = subprocess.run([sys.executable, "-c", probe], capture_output=True, text=True, timeout=60) + assert proc.returncode == 0, proc.stderr + assert json.loads(proc.stdout.strip().splitlines()[-1]) == [] From b0aea7bbafae52edb893c9bdbb92f1a82f544096 Mon Sep 17 00:00:00 2001 From: Michael Spilman <130705596+TheDrOnos@users.noreply.github.com> Date: Fri, 25 Sep 2026 09:14:14 -0500 Subject: [PATCH 07/13] remote_client: connect/send/recv/close over a blocking socket; recv waits on a selector --- de_shell/remote_client.py | 102 +++++++++++++++++++- tests/test_boundary.py | 1 + tests/test_remote_client.py | 180 +++++++++++++++++++++++++++++++++++- 3 files changed, 279 insertions(+), 4 deletions(-) diff --git a/de_shell/remote_client.py b/de_shell/remote_client.py index d811746..5977550 100644 --- a/de_shell/remote_client.py +++ b/de_shell/remote_client.py @@ -4,19 +4,36 @@ The relay speaks the backend's own framing to a remote peer: the peer sends bare JSON lines; the relay sends ``PLOTAPP:\\n`` messages and ``PLOTBIN::\\n
`` binary frames. This module -decodes that stream. +is that peer: one blocking socket, no threads, no reconnect, no auth, no +policy. Those belong to the app that uses it. Standard library only: importable without numpy, anyplotlib, asyncio or Electron (tests/test_remote_client.py checks it in a clean interpreter). + +Errors are the standard library's: + +- ``TimeoutError``: ``recv(timeout)`` passed without a complete unit. Bytes + already received are kept for the next call. +- ``ConnectionError`` (or a subclass): the relay closed or reset the + connection, or this Connection was closed. +- ``ValueError``: a malformed ``PLOTBIN`` prefix or header, or bad JSON after + ``PLOTAPP:``. The unit is consumed first, so the next ``recv`` continues + with what follows. The offending text or bytes are on ``.line``. """ from __future__ import annotations import json import re +import selectors +import socket +import time + +__all__ = ["Connection", "connect"] _PLOTAPP = "PLOTAPP:" _PLOTBIN = b"PLOTBIN:" -_PREFIX = re.compile(rb"PLOTBIN:(\d+):(\d+)") +_PREFIX = re.compile(rb"PLOTBIN:(\d{1,15}):(\d{1,15})") # bounded: int() refuses > 4300 digits +_RECV_BYTES = 1 << 20 def _malformed(what: str, line: str | bytes) -> ValueError: @@ -73,7 +90,8 @@ def pop(self) -> tuple | None: if len(buf) < end: return None raw_header = bytes(buf[nl + 1:nl + 1 + hlen]) - payload = bytes(buf[nl + 1 + hlen:end]) + with memoryview(buf) as mv: # one copy; released before the resize below + payload = bytes(mv[nl + 1 + hlen:end]) self._consume(end) try: header = json.loads(raw_header.decode("utf-8")) @@ -101,3 +119,81 @@ def pop(self) -> tuple | None: def _consume(self, n: int) -> None: del self._buf[:n] self._scan = 0 + + +class Connection: + """One connection to a relay. Make it with :func:`connect`. + + ``recv`` may run on one thread while ``send`` runs on another: ``recv`` + waits on a selector with its own timeout and never changes the socket's + timeout, which stays the one ``connect`` set and bounds a ``send``. Two + concurrent ``recv`` calls, or two concurrent ``send`` calls, are not + supported. Call ``close`` once the reading thread has stopped. + """ + + def __init__(self, sock: socket.socket, timeout: float) -> None: + self._sock = sock + self._sock.settimeout(timeout) + self._selector = selectors.DefaultSelector() + self._selector.register(sock, selectors.EVENT_READ) + self._decoder = _Decoder() + self._closed = False + host, port = sock.getpeername()[:2] + self.remote: tuple[str, int] = (host, port) + + def send(self, obj: dict) -> None: + """One bare JSON line, sent whole. A non-finite number raises ValueError, sending nothing.""" + if self._closed: + raise ConnectionError("connection is closed") + line = json.dumps(obj, allow_nan=False, separators=(",", ":")) + "\n" + self._sock.sendall(line.encode("utf-8")) + + def recv(self, timeout: float | None = None) -> tuple: + """The next unit: ``("message", dict)``, ``("binary", header, payload)`` + or ``("stream", str)``. ``timeout=None`` waits indefinitely.""" + if self._closed: + raise ConnectionError("connection is closed") + deadline = None if timeout is None else time.monotonic() + timeout + while True: + unit = self._decoder.pop() + if unit is not None: + return unit + if deadline is None: + ready = self._selector.select(None) + else: + ready = self._selector.select(max(0.0, deadline - time.monotonic())) + if not ready: + raise TimeoutError( + f"no complete unit within {timeout} s " + f"({self._decoder.buffered} bytes of one buffered)") + chunk = self._sock.recv(_RECV_BYTES) + if not chunk: + if self._decoder.buffered: + raise ConnectionError( + f"the relay closed the connection inside an incomplete unit " + f"({self._decoder.buffered} bytes dropped)") + raise ConnectionError("the relay closed the connection") + self._decoder.feed(chunk) + + def close(self) -> None: + """Idempotent. Later ``send``/``recv`` raise ConnectionError.""" + if self._closed: + return + self._closed = True + self._selector.close() + try: + self._sock.shutdown(socket.SHUT_RDWR) + except OSError: + pass # already reset or never fully connected + self._sock.close() + + +def connect(host: str, port: int, timeout: float = 10.0) -> Connection: + """Connect to a relay. ``timeout`` bounds the connect and every ``send``.""" + sock = socket.create_connection((host, port), timeout=timeout) + try: + sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1) + return Connection(sock, timeout) + except BaseException: + sock.close() + raise diff --git a/tests/test_boundary.py b/tests/test_boundary.py index 34f4da8..75862da 100644 --- a/tests/test_boundary.py +++ b/tests/test_boundary.py @@ -32,6 +32,7 @@ "de_shell.debug_flags", "de_shell.compute", "de_shell.session", + "de_shell.remote_client", ] #: Absent from sys.modules after the imports above. diff --git a/tests/test_remote_client.py b/tests/test_remote_client.py index 4d2ce24..636d68b 100644 --- a/tests/test_remote_client.py +++ b/tests/test_remote_client.py @@ -12,11 +12,13 @@ import socket import subprocess import sys +import threading +import time import pytest from anyplotlib._binary_frame import encode_frame -from de_shell import ipc +from de_shell import ipc, remote_client from de_shell.remote_client import _Decoder NASTY = b"\nPLOTBIN:9:9\nPLOTAPP:{}\n" + bytes(i & 0xFF for i in range(3000)) @@ -129,6 +131,15 @@ def test_a_malformed_header_raises_after_consuming_the_frame(): assert decoder.pop() == ("message", {"type": "after"}) +def test_an_overlong_prefix_number_raises_with_the_line_and_the_decoder_recovers(): + decoder = _Decoder() + decoder.feed(b"PLOTBIN:" + b"9" * 5000 + b":1\nPLOTAPP:{}\n") + with pytest.raises(ValueError, match="PLOTBIN prefix") as excinfo: + decoder.pop() + assert excinfo.value.line == b"PLOTBIN:" + b"9" * 5000 + b":1" + assert decoder.pop() == ("message", {}) + + def test_crlf_is_stripped_and_bytes_that_are_not_utf8_become_replacement_characters(): decoder = _Decoder() decoder.feed(b"log line\r\n" + b"f\xff\xfe\n") @@ -147,3 +158,170 @@ def test_the_client_imports_only_the_standard_library(): proc = subprocess.run([sys.executable, "-c", probe], capture_output=True, text=True, timeout=60) assert proc.returncode == 0, proc.stderr assert json.loads(proc.stdout.strip().splitlines()[-1]) == [] + + +# --- Connection, on a loopback TCP pair --------------------------------------- + + +@pytest.fixture +def pair(): + """A client Connection and the raw server-side socket it is connected to.""" + server = socket.create_server(("127.0.0.1", 0)) + port = server.getsockname()[1] + conn = remote_client.connect("127.0.0.1", port, timeout=5.0) + peer, _ = server.accept() + server.close() + peer.settimeout(5.0) + try: + yield conn, peer, port + finally: + conn.close() + peer.close() + + +def _read_line(sock: socket.socket) -> bytes: + data = bytearray() + while not data.endswith(b"\n"): + chunk = sock.recv(65536) + assert chunk, "the client closed before a full line" + data += chunk + return bytes(data) + + +def test_connect_reports_the_remote(pair): + conn, _, port = pair + assert conn.remote == ("127.0.0.1", port) + + +def test_send_writes_one_compact_json_line(pair): + conn, peer, _ = pair + obj = {"type": "hello", "note": "εxx Å", "text": "a\nb"} + conn.send(obj) + line = _read_line(peer) + assert line == (json.dumps(obj, separators=(",", ":")) + "\n").encode("ascii") + assert line.count(b"\n") == 1 + assert json.loads(line) == obj + + +def test_send_refuses_a_non_finite_number_and_writes_nothing(pair): + conn, peer, _ = pair + with pytest.raises(ValueError): + conn.send({"type": "fit", "value": float("nan")}) + conn.send({"ok": True}) + assert _read_line(peer) == b'{"ok":true}\n' + + +def test_recv_returns_each_unit_kind(pair): + conn, peer, _ = pair + peer.sendall( + b'PLOTAPP:{"type":"a"}\n' + + encode_frame("f", "k", {}, b"\x00\n\x01") + + b"log line\r\n" + ) + assert conn.recv(timeout=5) == ("message", {"type": "a"}) + assert conn.recv(timeout=5) == ("binary", {"fig_id": "f", "key": "k"}, b"\x00\n\x01") + assert conn.recv(timeout=5) == ("stream", "log line") + + +def test_recv_timeout_raises_and_keeps_the_partial_unit(pair): + conn, peer, _ = pair + frame = encode_frame("f", "image", {"n": 1}, b"\x00\n\x01PLOTBIN:") + peer.sendall(b'PLOTAPP:{"type":"par') + with pytest.raises(TimeoutError): + conn.recv(timeout=0.2) + peer.sendall(b'tial"}\n' + frame[:10]) + assert conn.recv(timeout=5) == ("message", {"type": "partial"}) + with pytest.raises(TimeoutError): + conn.recv(timeout=0.2) + peer.sendall(frame[10:]) + assert conn.recv(timeout=5) == ( + "binary", {"n": 1, "fig_id": "f", "key": "image"}, b"\x00\n\x01PLOTBIN:", + ) + + +def test_units_received_before_eof_come_first_then_connection_error(pair): + conn, peer, _ = pair + peer.sendall(b'PLOTAPP:{"type":"last"}\n') + peer.shutdown(socket.SHUT_WR) + assert conn.recv(timeout=5) == ("message", {"type": "last"}) + with pytest.raises(ConnectionError): + conn.recv(timeout=5) + + +def test_eof_inside_a_unit_raises_connection_error(pair): + conn, peer, _ = pair + peer.sendall(b'PLOTAPP:{"type":') + peer.shutdown(socket.SHUT_WR) + with pytest.raises(ConnectionError, match="incomplete"): + conn.recv(timeout=5) + + +def test_a_malformed_unit_on_the_wire_raises_value_error_and_the_next_one_decodes(pair): + conn, peer, _ = pair + peer.sendall(b'PLOTBIN:12:x\nPLOTAPP:{"type":"after"}\n') + with pytest.raises(ValueError) as excinfo: + conn.recv(timeout=5) + assert excinfo.value.line == b"PLOTBIN:12:x" + assert conn.recv(timeout=5) == ("message", {"type": "after"}) + + +def test_close_is_idempotent_and_later_calls_raise_connection_error(pair): + conn, _, _ = pair + conn.close() + conn.close() + with pytest.raises(ConnectionError): + conn.recv(timeout=0.1) + with pytest.raises(ConnectionError): + conn.send({"type": "late"}) + + +def test_a_short_recv_poll_on_one_thread_does_not_break_a_long_send_on_another(pair): + # The connector's shape: one thread polls recv() with a short timeout while + # another sends. A per-call socket timeout would cut a blocked send short + # (TimeoutError partway through a line: a corrupt stream); the client waits + # on a selector instead. Many 1 MiB sends, not one big one: Windows accepts + # a single large non-blocking send whole, so only later sends ever block. + conn, peer, _ = pair + count = 32 + stop = threading.Event() + errors: list[BaseException] = [] + + def poll() -> None: + while not stop.is_set(): + try: + conn.recv(timeout=0.01) + except TimeoutError: + continue + except BaseException as e: + errors.append(e) + return + + received = bytearray() + + def slow_reader() -> None: + time.sleep(0.3) # the sends block on full socket buffers while the poller runs + try: + while received.count(b"\n") < count: + chunk = peer.recv(1 << 16) + if not chunk: + return + received.extend(chunk) + except OSError: + return # the sender gave up; the assertions below say why + + poller = threading.Thread(target=poll, daemon=True) + reader = threading.Thread(target=slow_reader, daemon=True) + poller.start() + reader.start() + block = "x" * (1 << 20) + try: + for i in range(count): + conn.send({"type": "blob", "i": i, "data": block}) + finally: + reader.join(timeout=30) + stop.set() + poller.join(timeout=5) + assert errors == [] + lines = bytes(received).splitlines() + assert [json.loads(line)["i"] for line in lines] == list(range(count)) + assert all(json.loads(line)["data"] == block for line in lines) From e241d2fd89848af2dcf148909e060728cfb81943 Mon Sep 17 00:00:00 2001 From: Michael Spilman <130705596+TheDrOnos@users.noreply.github.com> Date: Fri, 25 Sep 2026 09:17:42 -0500 Subject: [PATCH 08/13] remote_client: round trip against a real createRelay under Node --- tests/relay_echo_server.mjs | 36 +++++++++++++++++ tests/test_remote_client.py | 79 +++++++++++++++++++++++++++++++++++++ 2 files changed, 115 insertions(+) create mode 100644 tests/relay_echo_server.mjs diff --git a/tests/relay_echo_server.mjs b/tests/relay_echo_server.mjs new file mode 100644 index 0000000..4aac2fc --- /dev/null +++ b/tests/relay_echo_server.mjs @@ -0,0 +1,36 @@ +// relay_echo_server.mjs — a real createRelay under Node, for the cross-runtime +// round trip in test_remote_client.py (not collected by pytest). +// +// Listens on 127.0.0.1:0 and prints "PORT ". Every line is echoed back as +// {type: 'echo', line}. The first line on a connection also admits it and is +// followed by one PLOTBIN frame (PROBE in the test). A line containing "bye" +// is echoed, then the app closes the connection. Exits when stdin ends. +import { createRelay } from '../de_shell/js/main/relay.ts' + +const payload = Buffer.alloc(256 * 1024) +for (let i = 0; i < payload.length; i++) payload[i] = i & 0xff +payload.write('\nPLOTBIN:9:9\n', 1000, 'latin1') + +const greeted = new WeakSet() +const relay = await createRelay({ + host: '127.0.0.1', + port: 0, + onConnection: () => {}, + onLine: (c, line) => { + c.writeMessage({ type: 'echo', line }) + if (!greeted.has(c)) { + greeted.add(c) + c.admit() + c.writeBinary({ fig_id: 'probe', key: 'image', label: 'εxx Å' }, payload) + } + if (line.includes('"bye"')) c.close() + }, + onClose: () => {}, +}) + +process.stdout.write(`PORT ${relay.port}\n`) +process.stdin.resume() +process.stdin.on('end', () => { + relay.close().then(() => process.exit(0)) +}) +setTimeout(() => process.exit(3), 60_000).unref() diff --git a/tests/test_remote_client.py b/tests/test_remote_client.py index 636d68b..5518c51 100644 --- a/tests/test_remote_client.py +++ b/tests/test_remote_client.py @@ -9,11 +9,13 @@ from __future__ import annotations import json +import shutil import socket import subprocess import sys import threading import time +from pathlib import Path import pytest from anyplotlib._binary_frame import encode_frame @@ -325,3 +327,80 @@ def slow_reader() -> None: lines = bytes(received).splitlines() assert [json.loads(line)["i"] for line in lines] == list(range(count)) assert all(json.loads(line)["data"] == block for line in lines) + + +# --- cross-runtime: a real createRelay under Node ----------------------------- + +REPO = Path(__file__).resolve().parents[1] +HELPER = Path(__file__).with_name("relay_echo_server.mjs") + +# The frame the helper sends after the first line: every byte value, plus +# newline and marker lookalikes inside the payload. +_probe = bytearray(i & 0xFF for i in range(256 * 1024)) +_probe[1000:1013] = b"\nPLOTBIN:9:9\n" +PROBE = bytes(_probe) + + +def _node_that_strips_types() -> str: + node = shutil.which("node") + if node is None: + pytest.skip("node not on PATH") + probe = subprocess.run( + [node, "-p", "Boolean(process.features.typescript)"], + capture_output=True, text=True, timeout=60, + ) + if probe.stdout.strip() != "true": + pytest.skip("node on PATH cannot strip TypeScript types (needs Node >= 22.18)") + return node + + +@pytest.fixture +def relay(): + """The helper running a real createRelay; yields (port, process).""" + node = _node_that_strips_types() + proc = subprocess.Popen( + [node, str(HELPER)], cwd=REPO, + stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, + ) + try: + first = proc.stdout.readline().decode("ascii", "replace").strip() + if not first.startswith("PORT "): + proc.kill() + _, err = proc.communicate(timeout=10) + pytest.fail(f"relay helper did not start: {first!r}\n{err.decode('utf-8', 'replace')}") + yield int(first.split()[1]), proc + finally: + if proc.poll() is None: + proc.stdin.close() + try: + proc.wait(timeout=10) + except subprocess.TimeoutExpired: + proc.kill() + proc.wait() + + +def test_round_trip_against_a_real_relay(relay): + port, proc = relay + conn = remote_client.connect("127.0.0.1", port, timeout=10.0) + try: + assert conn.remote == ("127.0.0.1", port) + conn.send({"type": "hello", "client": "pytest"}) + assert conn.recv(timeout=10) == ( + "message", {"type": "echo", "line": '{"type":"hello","client":"pytest"}'}, + ) + kind, header, payload = conn.recv(timeout=10) + assert kind == "binary" + assert header == {"fig_id": "probe", "key": "image", "label": "εxx Å"} + assert payload == PROBE + action = {"type": "action", "name": "snap", "note": "εxx Å"} + conn.send(action) + kind, msg = conn.recv(timeout=10) + assert (kind, msg["type"], json.loads(msg["line"])) == ("message", "echo", action) + conn.send({"type": "bye"}) + assert conn.recv(timeout=10) == ("message", {"type": "echo", "line": '{"type":"bye"}'}) + with pytest.raises(ConnectionError): + conn.recv(timeout=10) # the app closed the connection after the echo + finally: + conn.close() + proc.stdin.close() + assert proc.wait(timeout=10) == 0 From 00797ecfc19ffffb728d4f51f2c0440cad2f44e0 Mon Sep 17 00:00:00 2001 From: Michael Spilman <130705596+TheDrOnos@users.noreply.github.com> Date: Fri, 25 Sep 2026 09:18:50 -0500 Subject: [PATCH 09/13] tests: close the relay helper's pipes in the round-trip fixture --- tests/test_remote_client.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/tests/test_remote_client.py b/tests/test_remote_client.py index 5518c51..5f24c46 100644 --- a/tests/test_remote_client.py +++ b/tests/test_remote_client.py @@ -377,6 +377,8 @@ def relay(): except subprocess.TimeoutExpired: proc.kill() proc.wait() + proc.stdout.close() + proc.stderr.close() def test_round_trip_against_a_real_relay(relay): From 9a85aa813934f7534af0983d8f616813a69a8132 Mon Sep 17 00:00:00 2001 From: Michael Spilman <130705596+TheDrOnos@users.noreply.github.com> Date: Fri, 25 Sep 2026 09:21:24 -0500 Subject: [PATCH 10/13] chore: news fragment for the relay --- upcoming_changes/+relay.new_feature.rst | 1 + 1 file changed, 1 insertion(+) create mode 100644 upcoming_changes/+relay.new_feature.rst diff --git a/upcoming_changes/+relay.new_feature.rst b/upcoming_changes/+relay.new_feature.rst new file mode 100644 index 0000000..615fff1 --- /dev/null +++ b/upcoming_changes/+relay.new_feature.rst @@ -0,0 +1 @@ +Apps gained a way to serve their backend's protocol to a remote client: ``createRelay`` (exported from ``@de/shell-main`` with the ``encodeMessage`` and ``encodeBinary`` encoders) listens on a TCP port and speaks the PLOTAPP/PLOTBIN framing, and ``de_shell.remote_client`` is its synchronous, standard-library Python client. From 6afaa131fe27d572f416d92f46ff4869ffd47bef Mon Sep 17 00:00:00 2001 From: Michael Spilman <130705596+TheDrOnos@users.noreply.github.com> Date: Fri, 25 Sep 2026 09:37:50 -0500 Subject: [PATCH 11/13] relay: refuse an empty host and type createServer by the one shape the relay calls --- de_shell/js/main/relay.test.ts | 22 +++++++++++++++++++++- de_shell/js/main/relay.ts | 7 ++++--- 2 files changed, 25 insertions(+), 4 deletions(-) diff --git a/de_shell/js/main/relay.test.ts b/de_shell/js/main/relay.test.ts index 24a0e4c..2cf83f4 100644 --- a/de_shell/js/main/relay.test.ts +++ b/de_shell/js/main/relay.test.ts @@ -123,7 +123,7 @@ function spyServer() { accepted.push({ socket, log }) }) return server - }) as typeof net.createServer + }) return { createServer, servers, accepted } } @@ -161,6 +161,26 @@ test('a bind to an address this machine lacks rejects with the OS error and leav assert.equal(spy.servers[0].listening, false) }) +test('an empty host rejects with a TypeError before any server is made', async () => { + const spy = spyServer() + await assert.rejects( + createRelay({ + host: '', + port: 0, + onConnection: () => {}, + onLine: () => {}, + onClose: () => {}, + createServer: spy.createServer, + }), + (err: unknown) => { + assert.ok(err instanceof TypeError, String(err)) + assert.match(err.message, /host is required/) + return true + }, + ) + assert.equal(spy.servers.length, 0, 'no server was made') +}) + test('a port already in use rejects with EADDRINUSE', async () => { const first = await startRelay() try { diff --git a/de_shell/js/main/relay.ts b/de_shell/js/main/relay.ts index ab414cd..bce9f6f 100644 --- a/de_shell/js/main/relay.ts +++ b/de_shell/js/main/relay.ts @@ -50,7 +50,7 @@ export interface Relay { } export interface RelayOptions { - /** Bound exactly: no fallback, no discovery. '0.0.0.0' binds every interface only when the app passes it. */ + /** Bound exactly: no fallback, no discovery. '0.0.0.0' binds every interface only when the app passes it; an empty host is refused. */ host: string /** 0 = ephemeral; the bound port is reported on the Relay. */ port: number @@ -59,7 +59,7 @@ export interface RelayOptions { onLine: (c: RelayConnection, line: string) => void /** Exactly once per connection; `err` is set for 'error'. */ onClose: (c: RelayConnection, reason: RelayCloseReason, err?: Error) => void - /** Accept to admit(); past it the connection is closed as 'hello-timeout'. Default 10 000. */ + /** Accept to admit(); past it the connection is closed as 'hello-timeout'. Must be > 0 (0 is not "disabled": it fires on the next tick). Default 10 000. */ helloTimeoutMs?: number /** Line cap in bytes (those before the '\n') until admit(). Default 64 KiB. */ preAdmitLineBytes?: number @@ -72,7 +72,7 @@ export interface RelayOptions { * factory returning tls.createServer(tlsOptions, connectionListener) fits. * Default net.createServer. */ - createServer?: typeof net.createServer + createServer?: (listener: (socket: net.Socket) => void) => net.Server /** The clock for the relay's timers; tests inject a fake one. Default: the global setTimeout. */ setTimeout?: (fn: () => void, ms: number) => unknown /** Pairs with setTimeout. Default: the global clearTimeout. */ @@ -89,6 +89,7 @@ function asError(e: unknown): Error { } export function createRelay(opts: RelayOptions): Promise { + if (!opts.host) return Promise.reject(new TypeError('createRelay: host is required; pass "0.0.0.0" to bind every interface')) const helloTimeoutMs = opts.helloTimeoutMs ?? 10_000 const preAdmitLineBytes = opts.preAdmitLineBytes ?? 64 * 1024 const lineBytes = opts.lineBytes ?? 16 * 1024 * 1024 From 635cc58cf9f86309d4fcd5a1266961606a566dcf Mon Sep 17 00:00:00 2001 From: Michael Spilman <130705596+TheDrOnos@users.noreply.github.com> Date: Fri, 25 Sep 2026 09:38:00 -0500 Subject: [PATCH 12/13] tests: arm the echo helper's watchdog before createRelay --- tests/relay_echo_server.mjs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/relay_echo_server.mjs b/tests/relay_echo_server.mjs index 4aac2fc..92869ff 100644 --- a/tests/relay_echo_server.mjs +++ b/tests/relay_echo_server.mjs @@ -11,6 +11,7 @@ const payload = Buffer.alloc(256 * 1024) for (let i = 0; i < payload.length; i++) payload[i] = i & 0xff payload.write('\nPLOTBIN:9:9\n', 1000, 'latin1') +setTimeout(() => process.exit(3), 60_000).unref() const greeted = new WeakSet() const relay = await createRelay({ host: '127.0.0.1', @@ -33,4 +34,3 @@ process.stdin.resume() process.stdin.on('end', () => { relay.close().then(() => process.exit(0)) }) -setTimeout(() => process.exit(3), 60_000).unref() From aaf84fcd5a67921a6e9fb33caebd8d7bcdb5b70e Mon Sep 17 00:00:00 2001 From: Michael Spilman <130705596+TheDrOnos@users.noreply.github.com> Date: Fri, 25 Sep 2026 11:12:14 -0500 Subject: [PATCH 13/13] chore: name the relay news fragment after its PR --- upcoming_changes/{+relay.new_feature.rst => 12.new_feature.rst} | 0 1 file changed, 0 insertions(+), 0 deletions(-) rename upcoming_changes/{+relay.new_feature.rst => 12.new_feature.rst} (100%) diff --git a/upcoming_changes/+relay.new_feature.rst b/upcoming_changes/12.new_feature.rst similarity index 100% rename from upcoming_changes/+relay.new_feature.rst rename to upcoming_changes/12.new_feature.rst