From c1d7d9130f4f587eca8afb2fda4f7f35c6bb29c9 Mon Sep 17 00:00:00 2001 From: Minggang Wang Date: Thu, 3 Sep 2026 15:44:00 +0800 Subject: [PATCH 1/8] [Web runtime] Add HTTP SSE action support --- bin/rclnodejs-web.js | 4 +- lib/runtime/transports/http.js | 184 +++++++++++++++++++++++++++++++- test/test-web-action.js | 125 +++++++++++++++++++++- web/client.js | 185 +++++++++++++++++++++++++++++++-- web/index.d.ts | 3 +- 5 files changed, 486 insertions(+), 15 deletions(-) diff --git a/bin/rclnodejs-web.js b/bin/rclnodejs-web.js index 3bfe1e616..df39da0b2 100755 --- a/bin/rclnodejs-web.js +++ b/bin/rclnodejs-web.js @@ -149,8 +149,8 @@ if (SUBCOMMANDS.has(argv[0])) { const httpHost = displayHost(cfg.http.host || cfg.host); const httpBase = cfg.http.basePath || cfg.path; const httpKinds = cfg.http.sse - ? 'call/publish + subscribe (SSE)' - : 'call/publish only'; + ? 'call/publish/action + subscribe (SSE)' + : 'call/publish/action'; process.stdout.write( ` also http://${httpHost}:${httpTransport.port}${httpBase} (${httpKinds})\n` ); diff --git a/lib/runtime/transports/http.js b/lib/runtime/transports/http.js index 4ef7c63de..e4d224282 100644 --- a/lib/runtime/transports/http.js +++ b/lib/runtime/transports/http.js @@ -19,6 +19,7 @@ const _STATUS_BY_CODE = Object.freeze({ not_exposed: 404, not_implemented: 501, unknown_kind: 400, + unknown_op: 400, invalid_frame: 400, invalid_json: 400, binary_unsupported: 400, @@ -30,6 +31,10 @@ const _STATUS_BY_CODE = Object.freeze({ call_failed: 500, publish_failed: 500, schema_violation: 400, + goal_rejected: 409, + action_failed: 500, + unknown_goal_id: 400, + missing_goal_id: 400, }); const _MAX_BODY_BYTES = 1 * 1024 * 1024; // 1 MiB cap on request bodies @@ -351,6 +356,165 @@ class HttpSseConnection extends Connection { } } +/** + * A long-lived {@link Connection} that streams a single action goal over + * Server-Sent Events (`text/event-stream`). + * + * One request → one goal. Drives the dispatcher with a single `action` + * `send_goal` frame, then relays `feedback` events and the terminal + * `result` event as SSE lines. Unlike {@link HttpSseConnection}, the + * stream always closes itself after the result — there is no ongoing + * subscription to keep open. + * + * No cancellation support: HTTP has no client→server back-channel once + * the request body is sent, so there is nowhere to carry a `cancel` frame. + * A goal that is sent runs to completion (or failure) server-side even if + * the client disconnects early. Use the WebSocket transport for + * cancelable actions. + */ +class HttpActionConnection extends Connection { + /** + * @param {import('http').IncomingMessage} req + * @param {import('http').ServerResponse} res + * @param {string} capability ROS action name to send the goal to. + * @param {*} payload The goal, parsed from the request body. + */ + constructor(req, res, capability, payload) { + super(); + this.req = req; + this.res = res; + this._capability = capability; + this._payload = payload; + this._streaming = false; + this._closed = false; + // Fixed goal id: one HTTP request carries exactly one goal. + this._goalId = 'http-action'; + + // Only `res` close signals a genuine client disconnect here. Unlike + // HttpSseConnection's GET requests (no body), this is a POST whose body + // was already fully read by `_readJsonBody` before this connection is + // even constructed — `req`'s own 'close' fires once that read completes, + // which is *not* a disconnect. Wiring it (as tried during development) + // ends the response prematurely, mid-goal, racing the still-in-flight + // rcl reply — the same hazard documented on HttpRequestConnection above. + const onClose = () => this._emitCloseOnce(); + res.on('close', onClose); + res.on('error', (e) => { + debug('action sse response error: %s', e.message); + this._emitCloseOnce(); + }); + } + + /** Kick off dispatch: emit a single `action` `send_goal` frame. */ + begin() { + this.emit('message', { + id: this._goalId, + kind: 'action', + op: 'send_goal', + capability: this._capability, + payload: this._payload, + }); + } + + send(frame) { + if (this._closed) return; + + // Feedback delivery — zero or more, only while streaming. + if (frame.event === 'feedback') { + this._ensureStream(); + this._writeEvent('feedback', frame.payload); + return; + } + + // Terminal result — exactly one, then the stream closes. + if (frame.event === 'result') { + this._ensureStream(); + if (frame.ok === false) { + this._writeEvent('error', { error: frame.error, code: frame.code }); + } else { + this._writeEvent('result', frame.payload, { status: frame.status }); + } + return this._emitCloseOnce(); + } + + // Goal accepted: switch the response into an event stream. + if (frame.ok === true) { + this._ensureStream(); + this._writeEvent('accepted', { capability: this._capability }); + return; + } + + // Goal rejected, or send_goal itself failed — headers not sent yet, + // surface as a normal HTTP error rather than an SSE stream. + const code = frame.code || 'internal_error'; + const status = _STATUS_BY_CODE[code] || 500; + this._writeError(status, code, frame.error || 'send_goal failed'); + this._emitCloseOnce(); + } + + /** Force-close the response. The dispatcher's `cleanup` will follow. */ + // eslint-disable-next-line no-unused-vars + close(code, reason) { + this._emitCloseOnce(); + } + + _ensureStream() { + if (this._streaming) return; + this._streaming = true; + this.res.writeHead(200, { + 'content-type': 'text/event-stream; charset=utf-8', + 'cache-control': 'no-cache, no-transform', + connection: 'keep-alive', + 'x-accel-buffering': 'no', + }); + if (typeof this.res.flushHeaders === 'function') { + this.res.flushHeaders(); + } + } + + _writeEvent(event, data, extra) { + if (this._closed) return; + const json = JSON.stringify( + extra ? { ...extra, payload: data } : (data ?? null) + ); + let chunk = `event: ${event}\n`; + for (const line of json.split('\n')) { + chunk += `data: ${line}\n`; + } + chunk += '\n'; + try { + this.res.write(chunk); + } catch (e) { + debug('action sse write failed: %s', e.message); + this._emitCloseOnce(); + } + } + + _writeError(status, code, message) { + const body = JSON.stringify({ ok: false, error: message, code }); + try { + this.res.writeHead(status, { + 'content-type': 'application/json; charset=utf-8', + 'content-length': Buffer.byteLength(body), + }); + this.res.end(body); + } catch (e) { + debug('action sse writeError failed: %s', e.message); + } + } + + _emitCloseOnce() { + if (this._closed) return; + this._closed = true; + try { + this.res.end(); + } catch { + /* response may already be finished */ + } + this.emit('close'); + } +} + /** * HTTP adapter for the Web Runtime. * @@ -359,6 +523,17 @@ class HttpSseConnection extends Connection { * POST /capability/call/ * POST /capability/publish/ * + * Also exposes `action` capabilities as a one-shot streaming request: + * + * POST /capability/action/ (text/event-stream response) + * + * The goal is sent as the JSON body; the response streams `accepted`, + * zero or more `feedback` events, then exactly one terminal `result` + * event before closing. There is no way to cancel an in-flight goal over + * HTTP (no back-channel once the request is sent) — use the WebSocket + * transport if cancellation is required. A client that disconnects + * early does not cancel the goal server-side; it runs to completion. + * * When `sse: true`, it additionally exposes `subscribe` as a Server-Sent * Events stream: * @@ -567,10 +742,10 @@ class HttpTransport extends TransportAdapter { return; } - if (kind !== 'call' && kind !== 'publish') { + if (kind !== 'call' && kind !== 'publish' && kind !== 'action') { return _writeJson(res, 404, { ok: false, - error: `unsupported kind over HTTP: ${kind} (use WebSocket for action)`, + error: `unsupported kind over HTTP: ${kind}`, code: 'unsupported_kind', }); } @@ -590,7 +765,10 @@ class HttpTransport extends TransportAdapter { const status = code === 'payload_too_large' ? 413 : 400; return _writeJson(res, status, { ok: false, error: err.message, code }); } - const conn = new HttpRequestConnection(req, res, kind, name, payload); + const conn = + kind === 'action' + ? new HttpActionConnection(req, res, name, payload) + : new HttpRequestConnection(req, res, kind, name, payload); try { this._onConnection(conn); conn.begin(); diff --git a/test/test-web-action.js b/test/test-web-action.js index 530ce2e2d..380addbe5 100644 --- a/test/test-web-action.js +++ b/test/test-web-action.js @@ -7,13 +7,19 @@ // http://www.apache.org/licenses/LICENSE-2.0 // Action capability dispatch coverage: raw wire-protocol frames (mirroring -// test-runtime.js) plus SDK-level (`rclnodejs/web`) WebSocket round-trips. +// test-runtime.js) plus SDK-level (`rclnodejs/web`) round-trips over both +// the WebSocket and HTTP transports. import assert from 'assert'; import { once } from 'node:events'; +import http from 'node:http'; import WebSocket, { WebSocketServer } from 'ws'; import rclnodejs from '../index.js'; -import { createRuntime, WebSocketTransport } from '../lib/runtime/index.js'; +import { + createRuntime, + WebSocketTransport, + HttpTransport, +} from '../lib/runtime/index.js'; import * as assertUtils from './utils.js'; // `web/` is ESM; dynamic import() defers loading the browser SDK until the test starts. @@ -31,6 +37,7 @@ describe('Action capability dispatch', function () { let runtime; let server; let wsUrl; + let httpUrl; function waitOpen(ws) { return new Promise((resolve, reject) => { @@ -92,11 +99,15 @@ describe('Action capability dispatch', function () { runtime = createRuntime({ node, - transport: new WebSocketTransport({ port: 0, host: '127.0.0.1' }), + transports: [ + new WebSocketTransport({ port: 0, host: '127.0.0.1' }), + new HttpTransport({ port: 0, host: '127.0.0.1' }), + ], }); runtime.expose({ action: { '/fibonacci': fibonacci } }); await runtime.start(); wsUrl = `ws://127.0.0.1:${runtime.transports[0].port}/capability`; + httpUrl = `http://127.0.0.1:${runtime.transports[1].port}`; }); after(async function () { @@ -492,5 +503,113 @@ describe('Action capability dispatch', function () { actionServer.destroy(); } }); + + it('sends a goal over HTTP (SSE) and awaits the result, with feedback', async function () { + const ros = await connect(httpUrl); + try { + const feedbacks = []; + const goal = await ros.action( + '/fibonacci', + { order: 5 }, + { onFeedback: (fb) => feedbacks.push(fb) } + ); + const result = await goal.result; + assert.deepStrictEqual(result.sequence, [1, 1, 2, 3]); + assert.strictEqual(goal.status, 'succeeded'); + assert.throws(() => { + goal.status = 'aborted'; + }, TypeError); + assert.strictEqual(feedbacks.length, 1); + assert.deepStrictEqual(feedbacks[0].sequence, [1, 1]); + } finally { + await ros.close(); + } + }); + + it('rejects cancel over HTTP with code:unsupported_kind', async function () { + const ros = await connect(httpUrl); + try { + const goal = await ros.action('/fibonacci', { order: 5 }); + await assert.rejects(goal.cancel(), (err) => { + assert.strictEqual(err.code, 'unsupported_kind'); + return true; + }); + await goal.result; + } finally { + await ros.close(); + } + }); + + it('rejects when an HTTP action stream ends without a result', async function () { + const truncatedServer = http.createServer((req, res) => { + res.writeHead(200, { 'content-type': 'text/event-stream' }); + res.end('event: accepted\ndata: {}\n\n'); + }); + await new Promise((resolve) => + truncatedServer.listen(0, '127.0.0.1', resolve) + ); + const address = truncatedServer.address(); + const ros = await connect(`http://127.0.0.1:${address.port}`); + try { + const goal = await ros.action('/fibonacci', { order: 5 }); + await assert.rejects(goal.result, (err) => { + assert.strictEqual(err.code, 'connection_lost'); + return true; + }); + } finally { + await ros.close(); + await new Promise((resolve, reject) => + truncatedServer.close((err) => (err ? reject(err) : resolve())) + ); + } + }); + + it('accepts CRLF-framed HTTP action events', async function () { + const crlfServer = http.createServer((req, res) => { + res.writeHead(200, { 'content-type': 'text/event-stream' }); + res.end( + 'event: result\r\ndata: {"payload":{"sequence":[1,2,3]}}\r\n\r\n' + ); + }); + await new Promise((resolve) => + crlfServer.listen(0, '127.0.0.1', resolve) + ); + const address = crlfServer.address(); + const ros = await connect(`http://127.0.0.1:${address.port}`); + try { + const goal = await ros.action('/fibonacci', { order: 5 }); + assert.deepStrictEqual(await goal.result, { sequence: [1, 2, 3] }); + assert.strictEqual(goal.status, 'unknown'); + } finally { + await ros.close(); + await new Promise((resolve, reject) => + crlfServer.close((err) => (err ? reject(err) : resolve())) + ); + } + }); + + it('rejects the result when HTTP stream setup fails', async function () { + const originalFetch = globalThis.fetch; + globalThis.fetch = async () => ({ + ok: true, + body: { + getReader() { + throw new Error('reader unavailable'); + }, + }, + }); + const ros = await connect('http://127.0.0.1:1'); + try { + const goal = await ros.action('/fibonacci', { order: 5 }); + await assert.rejects(goal.result, (err) => { + assert.strictEqual(err.code, 'network_error'); + assert.match(err.message, /reader unavailable/); + return true; + }); + } finally { + globalThis.fetch = originalFetch; + await ros.close(); + } + }); }); }); diff --git a/web/client.js b/web/client.js index e48b1f2b8..8c302ec95 100644 --- a/web/client.js +++ b/web/client.js @@ -413,11 +413,7 @@ class _WsLink { const goal = this._goals.get(frame.goalId); this._goals.delete(frame.goalId); if (goal) { - goal.setStatus( - ['succeeded', 'canceled', 'aborted'].includes(frame.status) - ? frame.status - : 'unknown' - ); + goal.setStatus(_normaliseActionStatus(frame.status)); if (frame.ok === false) { goal.rejectResult( Object.assign(new Error(frame.error || 'action failed'), { @@ -493,6 +489,71 @@ class _HttpLink { return this._fetch('publish', capability, payload, /* expectBody */ false); } + /** + * Send an action goal over HTTP. The response streams `feedback` + * events (relayed to `onFeedback`) and one terminal `result` event. + * No cancellation support over HTTP — the returned handle's `cancel()` + * always rejects with `code: 'unsupported_kind'`. Use the WebSocket + * transport for cancelable actions. + */ + async action(capability, payload, { onFeedback } = {}) { + const url = this.baseUrl + '/action/' + _encodeRosName(capability); + let res; + try { + res = await fetch(url, { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify(payload ?? {}), + }); + } catch (e) { + throw Object.assign(new Error(`HTTP request failed: ${e.message}`), { + code: 'network_error', + }); + } + + if (!res.ok || !res.body) { + let err = {}; + try { + err = await res.json(); + } catch (_) { + // non-JSON error body; fall back to the generic message below + } + throw Object.assign(new Error(err.error || `HTTP ${res.status}`), { + code: err.code || 'http_' + res.status, + status: res.status, + }); + } + + let resolveResult, rejectResult; + const result = new Promise((res2, rej2) => { + resolveResult = res2; + rejectResult = rej2; + }); + let status; + _pumpActionStream( + res.body, + onFeedback, + resolveResult, + rejectResult, + (value) => { + status = value; + } + ); + return { + goalId: 'http-action', + result, + get status() { + return status; + }, + cancel: () => + Promise.reject( + Object.assign(new Error('action cancel is not supported over HTTP'), { + code: 'unsupported_kind', + }) + ), + }; + } + async _fetch(kind, capability, payload, expectBody) { const url = this.baseUrl + '/' + kind + '/' + _encodeRosName(capability); let res; @@ -538,6 +599,12 @@ class _HttpLink { } } +function _normaliseActionStatus(status) { + return ['succeeded', 'canceled', 'aborted'].includes(status) + ? status + : 'unknown'; +} + function _connectionLostError(reconnecting = true) { return Object.assign( new Error( @@ -547,6 +614,107 @@ function _connectionLostError(reconnecting = true) { ); } +/** + * Read an SSE response body (from an action `fetch()`), relaying + * `feedback` events to `onFeedback` and settling `result`/`error` events + * against the goal's result promise. Runs detached from the caller’s + * await chain — the returned handle's `result` promise is what the + * caller actually awaits. + */ +async function _pumpActionStream( + body, + onFeedback, + resolveResult, + rejectResult, + setStatus +) { + let reader; + let buffer = ''; + let terminalReceived = false; + try { + reader = body.getReader(); + const decoder = new TextDecoder(); + for (;;) { + const { done, value } = await reader.read(); + if (done) break; + buffer += decoder.decode(value, { stream: true }); + let boundary; + while ((boundary = /\r?\n\r?\n/.exec(buffer)) !== null) { + const chunk = buffer.slice(0, boundary.index); + buffer = buffer.slice(boundary.index + boundary[0].length); + const { event, data } = _parseSseChunk(chunk); + if (event === 'feedback' && onFeedback) { + try { + onFeedback(data); + } catch (_) { + // user callback errors don't break the stream + } + } else if (event === 'result') { + terminalReceived = true; + setStatus(_normaliseActionStatus(data?.status)); + resolveResult( + data && data.payload !== undefined ? data.payload : data + ); + return; + } else if (event === 'error') { + terminalReceived = true; + setStatus(_normaliseActionStatus(data?.status)); + rejectResult( + Object.assign(new Error((data && data.error) || 'action failed'), { + code: data && data.code, + }) + ); + return; + } + } + } + rejectResult( + Object.assign(new Error('action stream ended before a terminal result'), { + code: 'connection_lost', + }) + ); + } catch (e) { + rejectResult( + Object.assign(new Error(`action stream read failed: ${e.message}`), { + code: 'network_error', + }) + ); + } finally { + if (reader && terminalReceived) { + try { + await reader.cancel(); + } catch { + /* stream may already be closed */ + } + } + if (reader) { + try { + reader.releaseLock(); + } catch { + /* stream may already be released */ + } + } + } +} + +/** Parse one `event:`/`data:` SSE block into `{event, data}`. */ +function _parseSseChunk(chunk) { + let event = 'message'; + const dataLines = []; + for (const line of chunk.split(/\r?\n/)) { + if (line.startsWith('event:')) event = line.slice(6).trim(); + else if (line.startsWith('data:')) + dataLines.push(line.slice(5).replace(/^ /, '')); + } + let data; + try { + data = JSON.parse(dataLines.join('\n')); + } catch (_) { + data = undefined; + } + return { event, data }; +} + // ROS names always start with `/`. Encode each path segment so that // names with `~`, `:`, etc. survive routing while keeping the leading // slash stripped (the server adds it back). @@ -749,7 +917,8 @@ export class RosClient { /** * Send an action goal. Returns `{ goalId, result, status, cancel() }` where * `result` is a Promise resolving with the action result, and `cancel()` - * requests cancellation over WebSocket. + * requests cancellation over WebSocket; for goals sent over HTTP, + * `cancel()` rejects with `code: 'unsupported_kind'`. * Read `status` after awaiting `result` to distinguish success, cancellation, * and abortion without changing the result payload. * @param {string} capability @@ -764,6 +933,10 @@ export class RosClient { 'action(capability, payload, options): onFeedback must be a function' ); } + if (this._closed) throw new Error('connection closed'); + if (this._http) { + return this._http.action(capability, payload, { onFeedback }); + } const ws = await this._ensureWs(); return ws.action(capability, payload, { onFeedback }); } diff --git a/web/index.d.ts b/web/index.d.ts index 3997d0776..cfc62002a 100644 --- a/web/index.d.ts +++ b/web/index.d.ts @@ -133,7 +133,8 @@ declare module 'rclnodejs/web' { /** * Handle for an in-flight action goal, returned by {@link RosClient.action}. * - * `cancel()` requests cancellation of the goal over WebSocket. + * `cancel()` requests cancellation over WebSocket. For goals sent over HTTP, + * it rejects with `code: 'unsupported_kind'`; closing the stream does not cancel the goal. * `result` resolves with the ROS payload even when canceled or aborted; * inspect `status` after awaiting it to determine the outcome. */ From c70a37cf8f110dc1fe6aebed5ab1c84645d9de34 Mon Sep 17 00:00:00 2001 From: Minggang Wang Date: Fri, 11 Sep 2026 10:52:24 +0800 Subject: [PATCH 2/8] Clarify comments --- lib/runtime/transports/http.js | 50 ++++++---------------------------- test/test-web-action.js | 4 +-- web/client.js | 30 +++++--------------- web/index.d.ts | 4 +-- 4 files changed, 19 insertions(+), 69 deletions(-) diff --git a/lib/runtime/transports/http.js b/lib/runtime/transports/http.js index e4d224282..a7dc0cb09 100644 --- a/lib/runtime/transports/http.js +++ b/lib/runtime/transports/http.js @@ -19,7 +19,6 @@ const _STATUS_BY_CODE = Object.freeze({ not_exposed: 404, not_implemented: 501, unknown_kind: 400, - unknown_op: 400, invalid_frame: 400, invalid_json: 400, binary_unsupported: 400, @@ -356,28 +355,13 @@ class HttpSseConnection extends Connection { } } -/** - * A long-lived {@link Connection} that streams a single action goal over - * Server-Sent Events (`text/event-stream`). - * - * One request → one goal. Drives the dispatcher with a single `action` - * `send_goal` frame, then relays `feedback` events and the terminal - * `result` event as SSE lines. Unlike {@link HttpSseConnection}, the - * stream always closes itself after the result — there is no ongoing - * subscription to keep open. - * - * No cancellation support: HTTP has no client→server back-channel once - * the request body is sent, so there is nowhere to carry a `cancel` frame. - * A goal that is sent runs to completion (or failure) server-side even if - * the client disconnects early. Use the WebSocket transport for - * cancelable actions. - */ +/** Streams one action over SSE; disconnecting does not cancel the goal. */ class HttpActionConnection extends Connection { /** * @param {import('http').IncomingMessage} req * @param {import('http').ServerResponse} res - * @param {string} capability ROS action name to send the goal to. - * @param {*} payload The goal, parsed from the request body. + * @param {string} capability ROS action name. + * @param {*} payload Parsed goal. */ constructor(req, res, capability, payload) { super(); @@ -387,16 +371,10 @@ class HttpActionConnection extends Connection { this._payload = payload; this._streaming = false; this._closed = false; - // Fixed goal id: one HTTP request carries exactly one goal. + // Each HTTP connection owns one goal, so a fixed ID is sufficient. this._goalId = 'http-action'; - // Only `res` close signals a genuine client disconnect here. Unlike - // HttpSseConnection's GET requests (no body), this is a POST whose body - // was already fully read by `_readJsonBody` before this connection is - // even constructed — `req`'s own 'close' fires once that read completes, - // which is *not* a disconnect. Wiring it (as tried during development) - // ends the response prematurely, mid-goal, racing the still-in-flight - // rcl reply — the same hazard documented on HttpRequestConnection above. + // Use res: req.close marks POST body completion, not client disconnect. const onClose = () => this._emitCloseOnce(); res.on('close', onClose); res.on('error', (e) => { @@ -405,7 +383,6 @@ class HttpActionConnection extends Connection { }); } - /** Kick off dispatch: emit a single `action` `send_goal` frame. */ begin() { this.emit('message', { id: this._goalId, @@ -419,14 +396,12 @@ class HttpActionConnection extends Connection { send(frame) { if (this._closed) return; - // Feedback delivery — zero or more, only while streaming. if (frame.event === 'feedback') { this._ensureStream(); this._writeEvent('feedback', frame.payload); return; } - // Terminal result — exactly one, then the stream closes. if (frame.event === 'result') { this._ensureStream(); if (frame.ok === false) { @@ -437,22 +412,19 @@ class HttpActionConnection extends Connection { return this._emitCloseOnce(); } - // Goal accepted: switch the response into an event stream. if (frame.ok === true) { this._ensureStream(); this._writeEvent('accepted', { capability: this._capability }); return; } - // Goal rejected, or send_goal itself failed — headers not sent yet, - // surface as a normal HTTP error rather than an SSE stream. + // Failures before streaming use HTTP errors, not SSE events. const code = frame.code || 'internal_error'; const status = _STATUS_BY_CODE[code] || 500; this._writeError(status, code, frame.error || 'send_goal failed'); this._emitCloseOnce(); } - /** Force-close the response. The dispatcher's `cleanup` will follow. */ // eslint-disable-next-line no-unused-vars close(code, reason) { this._emitCloseOnce(); @@ -523,16 +495,12 @@ class HttpActionConnection extends Connection { * POST /capability/call/ * POST /capability/publish/ * - * Also exposes `action` capabilities as a one-shot streaming request: + * Action goals use a JSON POST with an SSE response: * * POST /capability/action/ (text/event-stream response) * - * The goal is sent as the JSON body; the response streams `accepted`, - * zero or more `feedback` events, then exactly one terminal `result` - * event before closing. There is no way to cancel an in-flight goal over - * HTTP (no back-channel once the request is sent) — use the WebSocket - * transport if cancellation is required. A client that disconnects - * early does not cancel the goal server-side; it runs to completion. + * Streams `accepted`, `feedback`, then one terminal `result` or `error`. + * Use WebSocket for cancellation; an HTTP disconnect does not cancel the goal. * * When `sse: true`, it additionally exposes `subscribe` as a Server-Sent * Events stream: diff --git a/test/test-web-action.js b/test/test-web-action.js index 380addbe5..daac1c0cd 100644 --- a/test/test-web-action.js +++ b/test/test-web-action.js @@ -6,9 +6,7 @@ // // http://www.apache.org/licenses/LICENSE-2.0 -// Action capability dispatch coverage: raw wire-protocol frames (mirroring -// test-runtime.js) plus SDK-level (`rclnodejs/web`) round-trips over both -// the WebSocket and HTTP transports. +// Action protocol and SDK tests over WebSocket and HTTP/SSE. import assert from 'assert'; import { once } from 'node:events'; diff --git a/web/client.js b/web/client.js index 8c302ec95..ca89d1fb8 100644 --- a/web/client.js +++ b/web/client.js @@ -489,13 +489,7 @@ class _HttpLink { return this._fetch('publish', capability, payload, /* expectBody */ false); } - /** - * Send an action goal over HTTP. The response streams `feedback` - * events (relayed to `onFeedback`) and one terminal `result` event. - * No cancellation support over HTTP — the returned handle's `cancel()` - * always rejects with `code: 'unsupported_kind'`. Use the WebSocket - * transport for cancelable actions. - */ + /** POST a goal and stream feedback/results; use WebSocket for cancellation. */ async action(capability, payload, { onFeedback } = {}) { const url = this.baseUrl + '/action/' + _encodeRosName(capability); let res; @@ -516,7 +510,7 @@ class _HttpLink { try { err = await res.json(); } catch (_) { - // non-JSON error body; fall back to the generic message below + // Fall back to the HTTP status for non-JSON errors. } throw Object.assign(new Error(err.error || `HTTP ${res.status}`), { code: err.code || 'http_' + res.status, @@ -614,13 +608,7 @@ function _connectionLostError(reconnecting = true) { ); } -/** - * Read an SSE response body (from an action `fetch()`), relaying - * `feedback` events to `onFeedback` and settling `result`/`error` events - * against the goal's result promise. Runs detached from the caller’s - * await chain — the returned handle's `result` promise is what the - * caller actually awaits. - */ +/** Runs independently; terminal SSE events settle the result promise. */ async function _pumpActionStream( body, onFeedback, @@ -647,7 +635,7 @@ async function _pumpActionStream( try { onFeedback(data); } catch (_) { - // user callback errors don't break the stream + // Callback errors must not interrupt result delivery. } } else if (event === 'result') { terminalReceived = true; @@ -697,7 +685,6 @@ async function _pumpActionStream( } } -/** Parse one `event:`/`data:` SSE block into `{event, data}`. */ function _parseSseChunk(chunk) { let event = 'message'; const dataLines = []; @@ -915,12 +902,9 @@ export class RosClient { } /** - * Send an action goal. Returns `{ goalId, result, status, cancel() }` where - * `result` is a Promise resolving with the action result, and `cancel()` - * requests cancellation over WebSocket; for goals sent over HTTP, - * `cancel()` rejects with `code: 'unsupported_kind'`. - * Read `status` after awaiting `result` to distinguish success, cancellation, - * and abortion without changing the result payload. + * Return `{ goalId, result, status, cancel() }` for an action goal. + * Await `result` for the ROS payload, then inspect terminal `status`. + * HTTP `cancel()` rejects with `unsupported_kind`; use WebSocket to cancel. * @param {string} capability * @param {*} payload The goal. * @param {object|null} [options] diff --git a/web/index.d.ts b/web/index.d.ts index cfc62002a..295007b9b 100644 --- a/web/index.d.ts +++ b/web/index.d.ts @@ -133,8 +133,8 @@ declare module 'rclnodejs/web' { /** * Handle for an in-flight action goal, returned by {@link RosClient.action}. * - * `cancel()` requests cancellation over WebSocket. For goals sent over HTTP, - * it rejects with `code: 'unsupported_kind'`; closing the stream does not cancel the goal. + * HTTP `cancel()` rejects with `unsupported_kind`; use WebSocket to cancel. + * Closing an HTTP stream does not cancel its ROS goal. * `result` resolves with the ROS payload even when canceled or aborted; * inspect `status` after awaiting it to determine the outcome. */ From 2cae0f57a6ebd2b06284d38871a34de777dcc239 Mon Sep 17 00:00:00 2001 From: Minggang Wang Date: Fri, 11 Sep 2026 11:05:33 +0800 Subject: [PATCH 3/8] Address comments --- lib/openapi.js | 96 ++++++++++++++++++++++++++++++++++++++--- test/test-openapi.js | 70 +++++++++++++++++++++++++++++- test/test-web-action.js | 44 +++++++++++++++++++ web/client.js | 35 ++++++++------- web/index.d.ts | 12 +++--- 5 files changed, 230 insertions(+), 27 deletions(-) diff --git a/lib/openapi.js b/lib/openapi.js index a0b3a2708..66fc137b4 100644 --- a/lib/openapi.js +++ b/lib/openapi.js @@ -172,8 +172,8 @@ function messageSchemaToJsonSchema(schema, components) { } /** - * Resolve a top-level capability type (message for publish/subscribe, - * service Request/Response for call) to a JSON Schema, without registering + * Resolve a capability payload (message, service Request/Response, or + * action Goal/Feedback/Result) to a JSON Schema, without registering * it as a component itself (the top-level request/response body is inlined * in the operation, only *nested* types become `$ref`d components — this * matches typical OpenAPI style for RPC-shaped APIs). @@ -195,8 +195,8 @@ function topLevelSchema(typeName, subType, components) { /** * Build a full OpenAPI 3.1 document from a capability registry snapshot - * (`CapabilityRegistry.list()`'s shape: `{call, publish, subscribe}`, each a - * `{name: typeName}` map). + * (`CapabilityRegistry.list()` has call, publish, subscribe, and action + * maps from capability names to ROS type names). * * No `servers` option: it's pure top-level metadata this function never * reads while building `paths`, so callers (e.g. the CLI's @@ -206,7 +206,7 @@ function topLevelSchema(typeName, subType, components) { * not the rclnodejs release that generated the document, and there's no * source for the former today — so it's a fixed `'0.0.0'` placeholder. * - * @param {{call: object, publish: object, subscribe: object}} capabilities + * @param {{call: object, publish: object, subscribe: object, action?: object}} capabilities * @param {object} [options] * @param {string} [options.title] * @param {string} [options.basePath] - default '/capability' @@ -305,6 +305,92 @@ function buildOpenApiDocument(capabilities, options = {}) { }; } + for (const [name, typeName] of Object.entries(capabilities.action || {})) { + const route = `${basePath}/action${name}`; + paths[route] = { + post: { + summary: `Send a ROS 2 action goal to ${name}`, + operationId: `action_${sanitizeName(name)}`, + 'x-ros-capability': { kind: 'action', name, type: typeName }, + description: + 'POST a JSON goal to receive accepted and feedback events, then a ' + + 'terminal result or error event. Result data contains status and payload. ' + + 'Use fetch() or curl -N; EventSource does not support POST. ' + + 'Use WebSocket for cancellation; disconnecting does not cancel the goal.', + requestBody: { + required: true, + content: { + 'application/json': { + schema: topLevelSchema(typeName, 'Goal', components), + }, + }, + }, + responses: { + 200: { + description: `Server-Sent Events stream for ${typeName}`, + content: { + 'text/event-stream': { + schema: { + description: + 'JSON data for each SSE event; titles name the event.', + anyOf: [ + { + title: 'accepted', + type: 'object', + properties: { capability: { type: 'string' } }, + required: ['capability'], + }, + { + title: 'feedback', + ...topLevelSchema(typeName, 'Feedback', components), + }, + { + title: 'result', + type: 'object', + properties: { + status: { + type: 'string', + enum: ['succeeded', 'canceled', 'aborted', 'unknown'], + }, + payload: topLevelSchema(typeName, 'Result', components), + }, + required: ['status', 'payload'], + }, + { + title: 'error', + type: 'object', + properties: { + error: { type: 'string' }, + code: { type: 'string' }, + }, + required: ['error', 'code'], + }, + ], + }, + }, + }, + }, + 404: notExposedResponse(), + 409: { + description: 'Action goal rejected before streaming', + content: { + 'application/json': { + schema: { + type: 'object', + properties: { + ok: { type: 'boolean', const: false }, + error: { type: 'string' }, + code: { type: 'string', const: 'goal_rejected' }, + }, + }, + }, + }, + }, + }, + }, + }; + } + return { openapi: '3.1.0', info: { title, version: '0.0.0' }, diff --git a/test/test-openapi.js b/test/test-openapi.js index 633720bf9..f1956437d 100644 --- a/test/test-openapi.js +++ b/test/test-openapi.js @@ -96,6 +96,7 @@ describe('lib/openapi.js', function () { call: { '/add_two_ints': 'example_interfaces/srv/AddTwoInts' }, publish: { '/chatter': 'std_msgs/msg/String' }, subscribe: { '/cmd_vel': 'geometry_msgs/msg/Twist' }, + action: { '/fibonacci': 'example_interfaces/action/Fibonacci' }, }; it('produces a valid-looking OpenAPI 3.1 document shell', function () { @@ -153,10 +154,72 @@ describe('lib/openapi.js', function () { ]); }); + it('documents an action capability as POST /capability/action/ over SSE', function () { + const doc = buildOpenApiDocument(capabilities); + const operation = doc.paths['/capability/action/fibonacci']?.post; + assert.ok(operation, 'expected an HTTP action operation'); + assert.strictEqual(operation.operationId, 'action_fibonacci'); + assert.deepStrictEqual(operation['x-ros-capability'], { + kind: 'action', + name: '/fibonacci', + type: 'example_interfaces/action/Fibonacci', + }); + assert.strictEqual(operation.requestBody.required, true); + const goalSchema = + operation.requestBody.content['application/json'].schema; + assert.deepStrictEqual(goalSchema.properties.order, { type: 'integer' }); + + const stream = operation.responses['200'].content['text/event-stream']; + assert.ok(stream, 'expected a text/event-stream response'); + const events = Object.fromEntries( + stream.schema.anyOf.map((schema) => [schema.title, schema]) + ); + assert.deepStrictEqual(Object.keys(events), [ + 'accepted', + 'feedback', + 'result', + 'error', + ]); + assert.strictEqual(events.accepted.properties.capability.type, 'string'); + const sequenceSchema = { type: 'array', items: { type: 'integer' } }; + assert.deepStrictEqual( + events.feedback.properties.sequence, + sequenceSchema + ); + assert.deepStrictEqual( + events.result.properties.payload.properties.sequence, + sequenceSchema + ); + assert.deepStrictEqual(events.result.properties.status.enum, [ + 'succeeded', + 'canceled', + 'aborted', + 'unknown', + ]); + assert.deepStrictEqual(events.result.required, ['status', 'payload']); + assert.strictEqual(events.error.properties.code.type, 'string'); + assert.strictEqual(events.error.properties.error.type, 'string'); + }); + + it('documents HTTP action errors before streaming', function () { + const doc = buildOpenApiDocument(capabilities); + const operation = doc.paths['/capability/action/fibonacci']?.post; + assert.ok(operation, 'expected an HTTP action operation'); + const rejected = + operation.responses['409'].content['application/json'].schema; + assert.strictEqual(rejected.properties.ok.const, false); + assert.strictEqual(rejected.properties.code.const, 'goal_rejected'); + const notExposed = + operation.responses['404'].content['application/json'].schema; + assert.strictEqual(notExposed.properties.code.const, 'not_exposed'); + }); + it('normalizes a trailing-slash basePath, matching HttpTransport, instead of emitting a double slash', function () { const doc = buildOpenApiDocument(capabilities, { basePath: '/api/' }); assert.ok(doc.paths['/api/call/add_two_ints']); assert.ok(!('/api//call/add_two_ints' in doc.paths)); + assert.ok(doc.paths['/api/action/fibonacci']); + assert.ok(!('/api//action/fibonacci' in doc.paths)); }); it('de-duplicates nested message types into one shared $ref component', function () { @@ -184,16 +247,21 @@ describe('lib/openapi.js', function () { describe('CLI subcommand (bin/rclnodejs-web.js openapi)', function () { this.timeout(10000); - it('openapi prints a document from --call/--publish flags, no ROS init needed', async function () { + it('openapi prints a document from --call/--action flags, no ROS init needed', async function () { const { code, stdout } = await runCli([ 'openapi', '--call', '/add_two_ints=example_interfaces/srv/AddTwoInts', + '--action', + '/fibonacci=example_interfaces/action/Fibonacci', ]); assert.strictEqual(code, 0); const doc = JSON.parse(stdout); assert.strictEqual(doc.openapi, '3.1.0'); assert.ok(doc.paths['/capability/call/add_two_ints']); + const action = doc.paths['/capability/action/fibonacci']?.post; + assert.ok(action, 'expected the CLI to include HTTP action capabilities'); + assert.ok(action.responses['200'].content['text/event-stream']); }); it('openapi routes use --path when http.basePath is not set, matching the server transport', async function () { diff --git a/test/test-web-action.js b/test/test-web-action.js index daac1c0cd..c7d4afeee 100644 --- a/test/test-web-action.js +++ b/test/test-web-action.js @@ -524,6 +524,23 @@ describe('Action capability dispatch', function () { } }); + it('assigns unique IDs to concurrent and sequential HTTP action handles', async function () { + const ros = await connect(httpUrl); + try { + const goals = await Promise.all([ + ros.action('/fibonacci', { order: 5 }), + ros.action('/fibonacci', { order: 5 }), + ]); + await Promise.all(goals.map((goal) => goal.result)); + const nextGoal = await ros.action('/fibonacci', { order: 5 }); + await nextGoal.result; + goals.push(nextGoal); + assert.strictEqual(new Set(goals.map((goal) => goal.goalId)).size, 3); + } finally { + await ros.close(); + } + }); + it('rejects cancel over HTTP with code:unsupported_kind', async function () { const ros = await connect(httpUrl); try { @@ -562,6 +579,33 @@ describe('Action capability dispatch', function () { } }); + for (const [description, data] of [ + ['invalid JSON', 'data: not-json\n'], + ['empty data', 'data:\n'], + ['missing data', ''], + ]) { + it(`rejects HTTP action results with ${description}`, async function () { + const resultServer = http.createServer((req, res) => { + res.writeHead(200, { 'content-type': 'text/event-stream' }); + res.end(`event: result\n${data}\n`); + }); + await new Promise((resolve) => + resultServer.listen(0, '127.0.0.1', resolve) + ); + const address = resultServer.address(); + const ros = await connect(`http://127.0.0.1:${address.port}`); + try { + const goal = await ros.action('/fibonacci', { order: 5 }); + await assert.rejects(goal.result, { code: 'invalid_response' }); + } finally { + await ros.close(); + await new Promise((resolve, reject) => + resultServer.close((err) => (err ? reject(err) : resolve())) + ); + } + }); + } + it('accepts CRLF-framed HTTP action events', async function () { const crlfServer = http.createServer((req, res) => { res.writeHead(200, { 'content-type': 'text/event-stream' }); diff --git a/web/client.js b/web/client.js index ca89d1fb8..200c49eba 100644 --- a/web/client.js +++ b/web/client.js @@ -24,8 +24,8 @@ // Two transports are supported, picked from the URL scheme: // // - ws:// / wss:// → WebSocket only (call/publish/subscribe/action). -// - http:// / https:// → HTTP for call/publish; subscribe and action -// lazily use a sibling WebSocket. +// - http:// / https:// → HTTP for call/publish; SSE for action. +// Subscribe lazily uses a sibling WebSocket. // - { http, ws } → explicit endpoint pair. let WS = globalThis.WebSocket; @@ -452,10 +452,8 @@ class _WsLink { } /** - * HTTP link. Speaks the L2 HTTP capability protocol used by - * `HttpTransport` on the server. Stateless — every `call`/`publish` - * is a one-shot `fetch()`. Does not support subscribe or actions; - * the public client uses the WebSocket link for those verbs. + * HTTP link for `call`/`publish` and SSE actions via `HttpTransport`. + * Subscriptions use the WebSocket link. */ class _HttpLink { constructor(baseUrl) { @@ -534,7 +532,7 @@ class _HttpLink { } ); return { - goalId: 'http-action', + goalId: _genId(), result, get status() { return status; @@ -639,6 +637,14 @@ async function _pumpActionStream( } } else if (event === 'result') { terminalReceived = true; + if (data === undefined) { + rejectResult( + Object.assign(new Error('invalid JSON in action result event'), { + code: 'invalid_response', + }) + ); + return; + } setStatus(_normaliseActionStatus(data?.status)); resolveResult( data && data.payload !== undefined ? data.payload : data @@ -719,9 +725,9 @@ function _encodeRosName(name) { * Picks a transport from the URL scheme: * * - `ws://`, `wss://` → WebSocket only (call/publish/subscribe/action). - * - `http://`, `https://` → HTTP for `call`/`publish`; `subscribe` and - * `action` lazily use a sibling WebSocket endpoint at the - * same host with `/capability` appended. + * - `http://`, `https://` → HTTP for `call`/`publish`; SSE for `action`. + * `subscribe` lazily uses a sibling WebSocket at the same host + * with `/capability` appended. * - object `{http, ws}` → both URLs spelled out explicitly. * * **Path conventions.** When a `ws://` / `wss://` URL is passed @@ -737,7 +743,7 @@ function _encodeRosName(name) { * // WebSocket-only (path defaults to /capability) * const ros = await connect('ws://robot.local:9000'); * - * // HTTP for call/publish, WS sibling for subscribe/action + * // HTTP for call/publish, SSE for action, WS sibling for subscribe * const ros = await connect('http://robot.local:9001'); * * // Split endpoints (e.g. WS behind a different proxy) @@ -771,7 +777,7 @@ export class RosClient { this._wsUrl = wsUrl; // Eagerly construct (but don't yet open) the WS link when the user // explicitly asked for it. When the WS URL was derived from an - // HTTP base, leave construction lazy until subscribe() or action(). + // HTTP base, leave construction lazy until subscribe(). this._ws = wsExplicit && wsUrl ? new _WsLink(wsUrl, this._wsOptions) : null; this._wsEager = !!wsExplicit; this._wsConnect = null; // memoised connect promise (in-flight or settled) @@ -816,8 +822,7 @@ export class RosClient { async connect() { if (this._closed) throw new Error('connection closed'); // Open HTTP eagerly (it's a no-op anyway). Defer a derived WebSocket - // until subscribe() or action(), allowing HTTP-only deployments to - // use call/publish without a WebSocket endpoint. + // until subscribe(), allowing call/publish/action without a WS endpoint. if (this._http) await this._http.connect(); if (this._wsEager) await this._ensureWs(); return this; @@ -957,7 +962,7 @@ function _resolveUrls(url) { } if (/^https?:\/\//i.test(url)) { // HTTP base URL: derive a sibling WS URL, but open it only for - // subscribe() or action(). + // subscribe(). return { httpUrl: url, wsUrl: _deriveWsSibling(url), wsExplicit: false }; } throw new TypeError( diff --git a/web/index.d.ts b/web/index.d.ts index 295007b9b..c2920642e 100644 --- a/web/index.d.ts +++ b/web/index.d.ts @@ -191,14 +191,14 @@ declare module 'rclnodejs/web' { /** * Browser-native Web Runtime client. * - * The user-facing verb API (`call` / `publish` / `subscribe`) is the - * same regardless of transport. The transport(s) used underneath - * are picked from the URL scheme passed to {@link connect}: + * The verb API (`call` / `publish` / `subscribe` / `action`) is the same + * regardless of transport. Transports are picked from the URL scheme + * passed to {@link connect}: * * - `ws://` / `wss://` — WebSocket only. - * - `http://` / `https://` — HTTP for `call`/`publish`; subscribe - * falls through to a sibling WebSocket endpoint at the same - * host with `/capability` appended. + * - `http://` / `https://` — HTTP for `call`/`publish`; SSE for `action`. + * `subscribe` falls through to a sibling WebSocket at the same host + * with `/capability` appended. * - {@link ConnectEndpoints} — both URLs spelled out. * * **Path conventions.** When a `ws://` / `wss://` URL is passed From b8d269235a77367cae5fc00b0e6d3c9fa7d7f51b Mon Sep 17 00:00:00 2001 From: Minggang Wang Date: Fri, 11 Sep 2026 11:23:51 +0800 Subject: [PATCH 4/8] Address comments --- lib/runtime/transports/http.js | 29 ++++- test/test-web-action.js | 230 ++++++++++++++++++++++++++++++++- web/client.js | 40 ++++-- 3 files changed, 285 insertions(+), 14 deletions(-) diff --git a/lib/runtime/transports/http.js b/lib/runtime/transports/http.js index a7dc0cb09..8c4eabbf5 100644 --- a/lib/runtime/transports/http.js +++ b/lib/runtime/transports/http.js @@ -362,13 +362,18 @@ class HttpActionConnection extends Connection { * @param {import('http').ServerResponse} res * @param {string} capability ROS action name. * @param {*} payload Parsed goal. + * @param {object} [options] + * @param {number} [options.keepAliveMs=15000] Heartbeat interval. */ - constructor(req, res, capability, payload) { + constructor(req, res, capability, payload, options = {}) { super(); this.req = req; this.res = res; this._capability = capability; this._payload = payload; + this._keepAliveMs = + options.keepAliveMs != null ? options.keepAliveMs : 15000; + this._keepAlive = null; this._streaming = false; this._closed = false; // Each HTTP connection owns one goal, so a fixed ID is sufficient. @@ -442,6 +447,20 @@ class HttpActionConnection extends Connection { if (typeof this.res.flushHeaders === 'function') { this.res.flushHeaders(); } + if (this._keepAliveMs > 0) { + this._keepAlive = setInterval(() => { + if (this._closed) return; + try { + this.res.write(': keep-alive\n\n'); + } catch (e) { + debug('action sse keep-alive write failed: %s', e.message); + this._emitCloseOnce(); + } + }, this._keepAliveMs); + if (typeof this._keepAlive.unref === 'function') { + this._keepAlive.unref(); + } + } } _writeEvent(event, data, extra) { @@ -478,6 +497,10 @@ class HttpActionConnection extends Connection { _emitCloseOnce() { if (this._closed) return; this._closed = true; + if (this._keepAlive) { + clearInterval(this._keepAlive); + this._keepAlive = null; + } try { this.res.end(); } catch { @@ -735,7 +758,9 @@ class HttpTransport extends TransportAdapter { } const conn = kind === 'action' - ? new HttpActionConnection(req, res, name, payload) + ? new HttpActionConnection(req, res, name, payload, { + keepAliveMs: this.sseKeepAliveMs, + }) : new HttpRequestConnection(req, res, kind, name, payload); try { this._onConnection(conn); diff --git a/test/test-web-action.js b/test/test-web-action.js index c7d4afeee..86bec75d3 100644 --- a/test/test-web-action.js +++ b/test/test-web-action.js @@ -9,8 +9,9 @@ // Action protocol and SDK tests over WebSocket and HTTP/SSE. import assert from 'assert'; -import { once } from 'node:events'; +import { EventEmitter, once } from 'node:events'; import http from 'node:http'; +import sinon from 'sinon'; import WebSocket, { WebSocketServer } from 'ws'; import rclnodejs from '../index.js'; import { @@ -555,6 +556,134 @@ describe('Action capability dispatch', function () { } }); + describe('HTTP action shutdown', function () { + let ros; + let originalFetch; + + beforeEach(async function () { + originalFetch = globalThis.fetch; + ros = await connect({ http: 'http://127.0.0.1:1' }); + }); + + afterEach(async function () { + globalThis.fetch = originalFetch; + await ros.close(); + }); + + it('aborts requests still waiting for response headers', async function () { + let signal; + let rejectFetch; + globalThis.fetch = (_url, options) => { + signal = options.signal; + return new Promise((resolve, reject) => { + rejectFetch = reject; + signal?.addEventListener('abort', () => reject(signal.reason), { + once: true, + }); + }); + }; + const outcome = ros + .action('/fibonacci', { order: 5 }) + .catch((error) => error); + try { + await ros.close(); + assert.ok(signal?.aborted, 'close must abort the pending fetch'); + assert.strictEqual((await outcome).code, 'connection_lost'); + } finally { + rejectFetch(new Error('test cleanup')); + await outcome; + } + }); + + it('aborts all idle streams and rejects their results', async function () { + const streams = []; + globalThis.fetch = async (_url, { signal }) => { + let controller; + const body = new ReadableStream({ + start(value) { + controller = value; + }, + }); + signal?.addEventListener( + 'abort', + () => controller.error(signal.reason), + { + once: true, + } + ); + streams.push({ signal, body, controller }); + return { ok: true, body }; + }; + const goals = await Promise.all([ + ros.action('/fibonacci', { order: 5 }), + ros.action('/fibonacci', { order: 5 }), + ]); + const outcomes = Promise.all( + goals.map((goal) => goal.result.catch((error) => error)) + ); + try { + await ros.close(); + assert.ok(streams.every(({ signal }) => signal?.aborted)); + for (const error of await outcomes) { + assert.strictEqual(error.code, 'connection_lost'); + } + await new Promise(setImmediate); + assert.ok(streams.every(({ body }) => !body.locked)); + assert.ok(goals.every((goal) => goal.status === undefined)); + } finally { + for (const { controller } of streams) { + controller.error(new Error('test cleanup')); + } + await outcomes; + } + }); + + it('stops buffered feedback when a callback closes the client', async function () { + let controller; + let cancelCount = 0; + const body = new ReadableStream({ + start(value) { + controller = value; + }, + cancel() { + cancelCount++; + }, + }); + globalThis.fetch = async () => ({ ok: true, body }); + let feedbackCount = 0; + let closing; + const goal = await ros.action( + '/fibonacci', + { order: 5 }, + { + onFeedback() { + feedbackCount++; + closing = ros.close(); + }, + } + ); + const outcome = goal.result.catch((error) => error); + try { + controller.enqueue( + new TextEncoder().encode( + 'event: feedback\ndata: {"sequence":[1]}\n\n' + + 'event: feedback\ndata: {"sequence":[1,1]}\n\n' + + 'event: result\ndata: {"status":"succeeded","payload":{}}\n\n' + ) + ); + const error = await outcome; + await closing; + assert.strictEqual(feedbackCount, 1); + assert.strictEqual(error.code, 'connection_lost'); + assert.strictEqual(goal.status, undefined); + assert.strictEqual(cancelCount, 1); + assert.strictEqual(body.locked, false); + } finally { + controller.error(new Error('test cleanup')); + } + }); + }); + it('rejects when an HTTP action stream ends without a result', async function () { const truncatedServer = http.createServer((req, res) => { res.writeHead(200, { 'content-type': 'text/event-stream' }); @@ -655,3 +784,102 @@ describe('Action capability dispatch', function () { }); }); }); + +describe('HTTP action heartbeats', function () { + let clock; + let connection; + let response; + + beforeEach(function () { + clock = sinon.useFakeTimers({ toFake: ['setInterval', 'clearInterval'] }); + }); + + afterEach(function () { + connection?.close(); + connection = null; + clock.restore(); + }); + + function openAction(options) { + const transport = new HttpTransport(options); + const request = Object.assign(new EventEmitter(), { + method: 'POST', + url: '/capability/action/fibonacci', + headers: { 'content-type': 'application/json' }, + }); + response = new EventEmitter(); + response.writeHead = sinon.stub().returns(response); + response.flushHeaders = sinon.spy(); + response.write = sinon.stub().returns(true); + response.end = sinon.spy(); + transport._onConnection = (value) => { + connection = value; + }; + transport._route(request, response); + request.emit('data', Buffer.from('{"order":5}')); + request.emit('end'); + assert.ok(connection); + } + + for (const [label, options, interval] of [ + ['default', {}, 15000], + ['configured', { sseKeepAliveMs: 10 }, 10], + ]) { + it(`uses the ${label} interval for actions without feedback`, function () { + openAction(options); + assert.strictEqual(clock.countTimers(), 0); + connection.send({ ok: true }); + clock.tick(interval - 1); + assert.ok(!response.write.calledWith(': keep-alive\n\n')); + clock.tick(1); + assert.ok(response.write.calledWith(': keep-alive\n\n')); + assert.strictEqual(clock.countTimers(), 1); + assert.strictEqual(connection._keepAlive.hasRef(), false); + }); + } + + it('disables heartbeats when sseKeepAliveMs is zero', function () { + openAction({ sseKeepAliveMs: 0 }); + connection.send({ ok: true }); + clock.tick(30000); + assert.strictEqual(clock.countTimers(), 0); + assert.ok(!response.write.calledWith(': keep-alive\n\n')); + }); + + for (const ending of ['result', 'error', 'disconnect', 'write failure']) { + it(`clears the timer after ${ending}`, function () { + openAction({ sseKeepAliveMs: 10 }); + connection.send({ ok: true }); + clock.tick(10); + assert.strictEqual(clock.countTimers(), 1); + if (ending === 'disconnect') { + response.emit('close'); + } else if (ending === 'write failure') { + response.write.throws(new Error('connection lost')); + clock.tick(10); + } else { + connection.send({ + event: 'result', + ok: ending === 'result', + status: 'succeeded', + payload: { sequence: [] }, + code: 'action_failed', + error: 'action failed', + }); + } + assert.strictEqual(clock.countTimers(), 0); + const writes = response.write.callCount; + clock.tick(30); + assert.strictEqual(response.write.callCount, writes); + }); + } + + it('does not start a timer for a rejected goal', function () { + openAction({ sseKeepAliveMs: 10 }); + connection.send({ ok: false, code: 'goal_rejected', error: 'rejected' }); + clock.tick(30); + assert.strictEqual(clock.countTimers(), 0); + assert.ok(response.writeHead.calledWith(409)); + assert.ok(response.write.notCalled); + }); +}); diff --git a/web/client.js b/web/client.js index 200c49eba..67909f915 100644 --- a/web/client.js +++ b/web/client.js @@ -467,6 +467,7 @@ class _HttpLink { this.baseUrl = trimmed.endsWith('/capability') ? trimmed : trimmed + '/capability'; + this._actionControllers = new Set(); } async connect() { @@ -476,7 +477,10 @@ class _HttpLink { } async close() { - // No-op for HTTP. + for (const controller of this._actionControllers) { + controller.abort(_connectionLostError(false)); + } + this._actionControllers.clear(); } call(capability, payload) { @@ -490,14 +494,21 @@ class _HttpLink { /** POST a goal and stream feedback/results; use WebSocket for cancellation. */ async action(capability, payload, { onFeedback } = {}) { const url = this.baseUrl + '/action/' + _encodeRosName(capability); + const controller = new AbortController(); + const { signal } = controller; + this._actionControllers.add(controller); let res; try { res = await fetch(url, { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify(payload ?? {}), + signal, }); + signal.throwIfAborted(); } catch (e) { + this._actionControllers.delete(controller); + if (signal.aborted) throw signal.reason; throw Object.assign(new Error(`HTTP request failed: ${e.message}`), { code: 'network_error', }); @@ -509,7 +520,10 @@ class _HttpLink { err = await res.json(); } catch (_) { // Fall back to the HTTP status for non-JSON errors. + } finally { + this._actionControllers.delete(controller); } + signal.throwIfAborted(); throw Object.assign(new Error(err.error || `HTTP ${res.status}`), { code: err.code || 'http_' + res.status, status: res.status, @@ -529,8 +543,9 @@ class _HttpLink { rejectResult, (value) => { status = value; - } - ); + }, + signal + ).finally(() => this._actionControllers.delete(controller)); return { goalId: _genId(), result, @@ -612,20 +627,23 @@ async function _pumpActionStream( onFeedback, resolveResult, rejectResult, - setStatus + setStatus, + signal ) { let reader; let buffer = ''; - let terminalReceived = false; try { reader = body.getReader(); const decoder = new TextDecoder(); for (;;) { + signal.throwIfAborted(); const { done, value } = await reader.read(); + signal.throwIfAborted(); if (done) break; buffer += decoder.decode(value, { stream: true }); let boundary; while ((boundary = /\r?\n\r?\n/.exec(buffer)) !== null) { + signal.throwIfAborted(); const chunk = buffer.slice(0, boundary.index); buffer = buffer.slice(boundary.index + boundary[0].length); const { event, data } = _parseSseChunk(chunk); @@ -636,7 +654,6 @@ async function _pumpActionStream( // Callback errors must not interrupt result delivery. } } else if (event === 'result') { - terminalReceived = true; if (data === undefined) { rejectResult( Object.assign(new Error('invalid JSON in action result event'), { @@ -651,7 +668,6 @@ async function _pumpActionStream( ); return; } else if (event === 'error') { - terminalReceived = true; setStatus(_normaliseActionStatus(data?.status)); rejectResult( Object.assign(new Error((data && data.error) || 'action failed'), { @@ -669,12 +685,14 @@ async function _pumpActionStream( ); } catch (e) { rejectResult( - Object.assign(new Error(`action stream read failed: ${e.message}`), { - code: 'network_error', - }) + signal.aborted + ? signal.reason + : Object.assign(new Error(`action stream read failed: ${e.message}`), { + code: 'network_error', + }) ); } finally { - if (reader && terminalReceived) { + if (reader) { try { await reader.cancel(); } catch { From f0d767b003484ed5dda88b1289f27aadf8ce2338 Mon Sep 17 00:00:00 2001 From: Minggang Wang Date: Fri, 11 Sep 2026 12:53:40 +0800 Subject: [PATCH 5/8] Address comments --- test/test-web-action.js | 15 +++++++++++++++ web/client.js | 7 ++++--- web/index.d.ts | 3 ++- 3 files changed, 21 insertions(+), 4 deletions(-) diff --git a/test/test-web-action.js b/test/test-web-action.js index 86bec75d3..4b4b6e05b 100644 --- a/test/test-web-action.js +++ b/test/test-web-action.js @@ -409,6 +409,21 @@ describe('Action capability dispatch', function () { } }); + it('cancels actions over WebSocket with explicit HTTP and WebSocket endpoints', async function () { + const ros = await connect({ http: httpUrl, ws: wsUrl }); + let result; + try { + const goal = await ros.action('/fibonacci', { order: 5 }); + result = goal.result.catch((error) => error); + assert.strictEqual(await goal.cancel(), undefined); + assert.deepStrictEqual(await result, { sequence: [] }); + assert.strictEqual(goal.status, 'canceled'); + } finally { + await ros.close(); + await result; + } + }); + for (const { name, terminal, errorCode } of [ { name: 'missing status', terminal: { payload: {} } }, { name: 'unknown status', terminal: { status: 'unknown', payload: {} } }, diff --git a/web/client.js b/web/client.js index 67909f915..354f5c009 100644 --- a/web/client.js +++ b/web/client.js @@ -26,7 +26,7 @@ // - ws:// / wss:// → WebSocket only (call/publish/subscribe/action). // - http:// / https:// → HTTP for call/publish; SSE for action. // Subscribe lazily uses a sibling WebSocket. -// - { http, ws } → explicit endpoint pair. +// - { http, ws } → HTTP for call/publish; WS for subscribe/action. let WS = globalThis.WebSocket; let _wsResolved = !!WS; @@ -746,7 +746,8 @@ function _encodeRosName(name) { * - `http://`, `https://` → HTTP for `call`/`publish`; SSE for `action`. * `subscribe` lazily uses a sibling WebSocket at the same host * with `/capability` appended. - * - object `{http, ws}` → both URLs spelled out explicitly. + * - object `{http, ws}` → HTTP for `call`/`publish`; + * WebSocket for `subscribe`/`action`. * * **Path conventions.** When a `ws://` / `wss://` URL is passed * without a path (or with just `/`), the SDK appends the runtime's @@ -941,7 +942,7 @@ export class RosClient { ); } if (this._closed) throw new Error('connection closed'); - if (this._http) { + if (this._http && !this._wsEager) { return this._http.action(capability, payload, { onFeedback }); } const ws = await this._ensureWs(); diff --git a/web/index.d.ts b/web/index.d.ts index c2920642e..0d5055585 100644 --- a/web/index.d.ts +++ b/web/index.d.ts @@ -199,7 +199,8 @@ declare module 'rclnodejs/web' { * - `http://` / `https://` — HTTP for `call`/`publish`; SSE for `action`. * `subscribe` falls through to a sibling WebSocket at the same host * with `/capability` appended. - * - {@link ConnectEndpoints} — both URLs spelled out. + * - {@link ConnectEndpoints} — when both URLs are provided, + * HTTP for `call`/`publish`; WebSocket for `subscribe`/`action`. * * **Path conventions.** When a `ws://` / `wss://` URL is passed * without a path (or with just `/`), the SDK appends the runtime's From 08e7ad7b9fa621b33f08ec186a71f6bceec05aaa Mon Sep 17 00:00:00 2001 From: Minggang Wang Date: Thu, 17 Sep 2026 13:26:54 +0800 Subject: [PATCH 6/8] Retry setting up ROS 2 --- .github/workflows/linux-arm64-build-and-test.yml | 9 +++++++++ .github/workflows/linux-x64-build-and-test.yml | 9 +++++++++ 2 files changed, 18 insertions(+) diff --git a/.github/workflows/linux-arm64-build-and-test.yml b/.github/workflows/linux-arm64-build-and-test.yml index 117bb4bbd..f961667db 100644 --- a/.github/workflows/linux-arm64-build-and-test.yml +++ b/.github/workflows/linux-arm64-build-and-test.yml @@ -60,7 +60,16 @@ jobs: - uses: actions/checkout@v7 - name: Setup ROS2 + id: setup_ros if: ${{ matrix.ros_distribution != 'rolling' && matrix.ros_distribution != 'lyrical' }} + continue-on-error: true + uses: ros-tooling/setup-ros@v0.7 + with: + required-ros-distributions: ${{ matrix.ros_distribution }} + + # Retry transient repository/download failures once; a second failure is fatal. + - name: Retry ROS2 setup + if: ${{ !cancelled() && steps.setup_ros.outcome == 'failure' }} uses: ros-tooling/setup-ros@v0.7 with: required-ros-distributions: ${{ matrix.ros_distribution }} diff --git a/.github/workflows/linux-x64-build-and-test.yml b/.github/workflows/linux-x64-build-and-test.yml index ffb351ece..27a5f7b5c 100644 --- a/.github/workflows/linux-x64-build-and-test.yml +++ b/.github/workflows/linux-x64-build-and-test.yml @@ -64,7 +64,16 @@ jobs: - uses: actions/checkout@v7 - name: Setup ROS2 + id: setup_ros if: ${{ matrix.ros_distribution != 'rolling' && matrix.ros_distribution != 'lyrical' }} + continue-on-error: true + uses: ros-tooling/setup-ros@v0.7 + with: + required-ros-distributions: ${{ matrix.ros_distribution }} + + # Retry transient repository/download failures once; a second failure is fatal. + - name: Retry ROS2 setup + if: ${{ !cancelled() && steps.setup_ros.outcome == 'failure' }} uses: ros-tooling/setup-ros@v0.7 with: required-ros-distributions: ${{ matrix.ros_distribution }} From fcf3d7d3c7276d16678191dcc580b8a6e73f8478 Mon Sep 17 00:00:00 2001 From: Minggang Wang Date: Thu, 17 Sep 2026 13:59:40 +0800 Subject: [PATCH 7/8] Address comments --- test/test-web-action.js | 69 ++++++++++++++++++++++++++++------------- web/client.js | 16 +++++++--- 2 files changed, 58 insertions(+), 27 deletions(-) diff --git a/test/test-web-action.js b/test/test-web-action.js index 4b4b6e05b..dbf3603db 100644 --- a/test/test-web-action.js +++ b/test/test-web-action.js @@ -750,29 +750,54 @@ describe('Action capability dispatch', function () { }); } - it('accepts CRLF-framed HTTP action events', async function () { - const crlfServer = http.createServer((req, res) => { - res.writeHead(200, { 'content-type': 'text/event-stream' }); - res.end( - 'event: result\r\ndata: {"payload":{"sequence":[1,2,3]}}\r\n\r\n' - ); - }); - await new Promise((resolve) => - crlfServer.listen(0, '127.0.0.1', resolve) - ); - const address = crlfServer.address(); - const ros = await connect(`http://127.0.0.1:${address.port}`); - try { - const goal = await ros.action('/fibonacci', { order: 5 }); - assert.deepStrictEqual(await goal.result, { sequence: [1, 2, 3] }); - assert.strictEqual(goal.status, 'unknown'); - } finally { - await ros.close(); - await new Promise((resolve, reject) => - crlfServer.close((err) => (err ? reject(err) : resolve())) - ); + for (const [label, lineEnding, separator] of [ + ['LF', '\n', '\n\n'], + ['CRLF', '\r\n', '\r\n\r\n'], + ['CR', '\r', '\r\r'], + ['mixed CR/LF/CRLF', '\r', '\r\n\n'], + ]) { + for (const byteByByte of [false, true]) { + it(`accepts ${label}-framed HTTP action events ${byteByByte ? 'one byte at a time' : 'in one chunk'}`, async function () { + const originalFetch = globalThis.fetch; + const ros = await connect('http://127.0.0.1:1'); + const feedbacks = []; + const stream = + `: keep-alive${separator}` + + `event: accepted${lineEnding}data: {"capability":"/fibonacci"}${separator}` + + `event: feedback${lineEnding}data: {"sequence":[1,1]}${separator}` + + `event: result${lineEnding}data: {"payload":${lineEnding}` + + `data: {"sequence":[1,2,3]}}${separator}`; + const bytes = new TextEncoder().encode(stream); + const body = new ReadableStream({ + start(controller) { + const chunkSize = byteByByte ? 1 : bytes.length; + for (let index = 0; index < bytes.length; index += chunkSize) { + controller.enqueue(bytes.subarray(index, index + chunkSize)); + if (byteByByte && bytes[index] === 13) { + controller.enqueue(new Uint8Array()); + } + } + controller.close(); + }, + }); + globalThis.fetch = async () => ({ ok: true, body }); + + try { + const goal = await ros.action( + '/fibonacci', + { order: 5 }, + { onFeedback: (feedback) => feedbacks.push(feedback) } + ); + assert.deepStrictEqual(await goal.result, { sequence: [1, 2, 3] }); + assert.strictEqual(goal.status, 'unknown'); + assert.deepStrictEqual(feedbacks, [{ sequence: [1, 1] }]); + } finally { + globalThis.fetch = originalFetch; + await ros.close(); + } + }); } - }); + } it('rejects the result when HTTP stream setup fails', async function () { const originalFetch = globalThis.fetch; diff --git a/web/client.js b/web/client.js index 354f5c009..f38807dc8 100644 --- a/web/client.js +++ b/web/client.js @@ -632,6 +632,7 @@ async function _pumpActionStream( ) { let reader; let buffer = ''; + let skipLeadingLF = false; try { reader = body.getReader(); const decoder = new TextDecoder(); @@ -640,12 +641,17 @@ async function _pumpActionStream( const { done, value } = await reader.read(); signal.throwIfAborted(); if (done) break; - buffer += decoder.decode(value, { stream: true }); + let text = decoder.decode(value, { stream: true }); + if (text.length === 0) continue; + // A CRLF split across reads is one line ending, not an empty line. + if (skipLeadingLF && text.startsWith('\n')) text = text.slice(1); + skipLeadingLF = text.endsWith('\r'); + buffer += text.replace(/\r\n?/g, '\n'); let boundary; - while ((boundary = /\r?\n\r?\n/.exec(buffer)) !== null) { + while ((boundary = buffer.indexOf('\n\n')) !== -1) { signal.throwIfAborted(); - const chunk = buffer.slice(0, boundary.index); - buffer = buffer.slice(boundary.index + boundary[0].length); + const chunk = buffer.slice(0, boundary); + buffer = buffer.slice(boundary + 2); const { event, data } = _parseSseChunk(chunk); if (event === 'feedback' && onFeedback) { try { @@ -712,7 +718,7 @@ async function _pumpActionStream( function _parseSseChunk(chunk) { let event = 'message'; const dataLines = []; - for (const line of chunk.split(/\r?\n/)) { + for (const line of chunk.split('\n')) { if (line.startsWith('event:')) event = line.slice(6).trim(); else if (line.startsWith('data:')) dataLines.push(line.slice(5).replace(/^ /, '')); From 853d5e6c1b8015aac4febcbd8c879304d840b368 Mon Sep 17 00:00:00 2001 From: Minggang Wang Date: Sun, 20 Sep 2026 13:02:04 +0800 Subject: [PATCH 8/8] Refactor SSE related classes --- lib/runtime/transports/http.js | 429 ++++++++++++++------------------- test/test-web-action.js | 251 ++++++++++++------- 2 files changed, 348 insertions(+), 332 deletions(-) diff --git a/lib/runtime/transports/http.js b/lib/runtime/transports/http.js index 8c4eabbf5..fdbf6fcb0 100644 --- a/lib/runtime/transports/http.js +++ b/lib/runtime/transports/http.js @@ -53,16 +53,16 @@ const _MAX_BODY_BYTES = 1 * 1024 * 1024; // 1 MiB cap on request bodies */ class HttpRequestConnection extends Connection { /** - * @param {import('http').IncomingMessage} req - * @param {import('http').ServerResponse} res + * @param {import('http').IncomingMessage} request + * @param {import('http').ServerResponse} response * @param {'call'|'publish'} kind * @param {string} capability * @param {*} payload */ - constructor(req, res, kind, capability, payload) { + constructor(request, response, kind, capability, payload) { super(); - this.req = req; - this.res = res; + this.request = request; + this.response = response; this._kind = kind; this._capability = capability; this._payload = payload; @@ -73,8 +73,8 @@ class HttpRequestConnection extends Connection { // and write the body directly, so the value just needs to be present. this._id = '__http__'; - // NOTE: we deliberately do *not* wire `req.on('close')` / - // `res.on('close')` to `_emitCloseOnce()` here. Doing so would + // NOTE: we deliberately do *not* wire `request.on('close')` / + // `response.on('close')` to `_emitCloseOnce()` here. Doing so would // catch client-side disconnects, but the close event also fires // on normal end-of-response, racing the still-in-flight rcl // reply callback and tearing down the dispatcher's lazy Client @@ -124,15 +124,15 @@ class HttpRequestConnection extends Connection { if (frame.ok === true) { if (this._kind === 'publish') { - this.res.writeHead(204).end(); + this.response.writeHead(204).end(); } else { // call: serialise the payload (may be undefined for void replies). const body = JSON.stringify(frame.payload ?? null); - this.res.writeHead(200, { + this.response.writeHead(200, { 'content-type': 'application/json; charset=utf-8', 'content-length': Buffer.byteLength(body), }); - this.res.end(body); + this.response.end(body); } return this._emitCloseOnce(); } @@ -166,113 +166,44 @@ class HttpRequestConnection extends Connection { _writeError(status, code, message) { const body = JSON.stringify({ ok: false, error: message, code }); try { - this.res.writeHead(status, { + this.response.writeHead(status, { 'content-type': 'application/json; charset=utf-8', 'content-length': Buffer.byteLength(body), }); - this.res.end(body); + this.response.end(body); } catch (e) { debug('writeError failed: %s', e.message); } } } -/** - * A long-lived {@link Connection} that streams a single subscription over - * Server-Sent Events (`text/event-stream`). - * - * Unlike {@link HttpRequestConnection} (one request → one reply), this - * connection stays open: it drives the dispatcher with exactly one - * `subscribe` frame and then relays every `{event:'message'}` delivery to - * the client as an SSE `message` event until the client disconnects. - * - * The dispatcher is unchanged — it sees the same `subscribe` frame and the - * same `{event:'message', subId, payload}` deliveries it would send over - * WebSocket. This class only reframes them as SSE lines. - * - * HTTP status semantics are preserved for the *establishment* of the - * stream: the `200 text/event-stream` headers are not written until the - * dispatcher acknowledges the subscribe (`{ok:true}`). If the subscribe is - * rejected (e.g. `not_exposed`), the headers have not been sent yet, so the - * failure surfaces as a normal JSON error response with the mapped status - * code — exactly like `call`/`publish`. - */ +/** Shared SSE framing, heartbeats, and response lifecycle. */ class HttpSseConnection extends Connection { /** - * @param {import('http').IncomingMessage} req - * @param {import('http').ServerResponse} res - * @param {string} capability ROS name to subscribe to. + * @param {import('http').IncomingMessage} request + * @param {import('http').ServerResponse} response * @param {object} [options] - * @param {number} [options.keepAliveMs=15000] Heartbeat comment interval. + * @param {number} [options.keepAliveMs=15000] Heartbeat comment interval. */ - constructor(req, res, capability, options = {}) { + constructor(request, response, options = {}) { super(); - this.req = req; - this.res = res; - this._capability = capability; + this.request = request; + this.response = response; this._keepAliveMs = options.keepAliveMs != null ? options.keepAliveMs : 15000; this._streaming = false; this._closed = false; this._keepAlive = null; - // Fixed subId: an SSE connection carries exactly one subscription. - this._subId = 'sse'; + // Use response: request.close marks POST body completion, not client disconnect. const onClose = () => this._emitCloseOnce(); - req.on('close', onClose); - res.on('close', onClose); - res.on('error', (e) => { + response.on('close', onClose); + response.on('error', (e) => { debug('sse response error: %s', e.message); this._emitCloseOnce(); }); } - /** - * Kick off dispatch: emit a single `subscribe` frame. The dispatcher - * replies synchronously with `{ok:true}` (→ start the stream) or an - * error (→ JSON failure response). - */ - begin() { - this.emit('message', { - id: this._subId, - kind: 'subscribe', - capability: this._capability, - }); - } - - send(frame) { - if (this._closed) return; - - // Subscription delivery — the steady-state case. - if (frame.event === 'message') { - this._ensureStream(); - this._writeEvent('message', frame.payload); - return; - } - - // Subscribe acknowledgement: switch the response into an event stream. - if (frame.ok === true) { - this._ensureStream(); - this._writeEvent('ready', { - capability: this._capability, - subId: this._subId, - }); - return; - } - - // Failure path. - const code = frame.code || 'internal_error'; - if (!this._streaming) { - // Headers not sent yet — surface as a normal HTTP error. - const status = _STATUS_BY_CODE[code] || 500; - this._writeError(status, code, frame.error || 'subscribe failed'); - } else { - // Already streaming — emit a terminal SSE error event, then close. - this._writeEvent('error', { error: frame.error, code }); - } - this._emitCloseOnce(); - } - /** Close the stream; the dispatcher's `cleanup` follows via `'close'`. */ // eslint-disable-next-line no-unused-vars close(code, reason) { @@ -282,21 +213,21 @@ class HttpSseConnection extends Connection { _ensureStream() { if (this._streaming) return; this._streaming = true; - this.res.writeHead(200, { + this.response.writeHead(200, { 'content-type': 'text/event-stream; charset=utf-8', 'cache-control': 'no-cache, no-transform', connection: 'keep-alive', // Disable proxy buffering (nginx) so events flush immediately. 'x-accel-buffering': 'no', }); - if (typeof this.res.flushHeaders === 'function') { - this.res.flushHeaders(); + if (typeof this.response.flushHeaders === 'function') { + this.response.flushHeaders(); } if (this._keepAliveMs > 0) { this._keepAlive = setInterval(() => { if (this._closed) return; try { - this.res.write(': keep-alive\n\n'); + this.response.write(': keep-alive\n\n'); } catch (e) { debug('sse keep-alive write failed: %s', e.message); this._emitCloseOnce(); @@ -319,7 +250,7 @@ class HttpSseConnection extends Connection { } chunk += '\n'; try { - this.res.write(chunk); + this.response.write(chunk); } catch (e) { debug('sse write failed: %s', e.message); this._emitCloseOnce(); @@ -329,11 +260,11 @@ class HttpSseConnection extends Connection { _writeError(status, code, message) { const body = JSON.stringify({ ok: false, error: message, code }); try { - this.res.writeHead(status, { + this.response.writeHead(status, { 'content-type': 'application/json; charset=utf-8', 'content-length': Buffer.byteLength(body), }); - this.res.end(body); + this.response.end(body); } catch (e) { debug('sse writeError failed: %s', e.message); } @@ -347,7 +278,7 @@ class HttpSseConnection extends Connection { this._keepAlive = null; } try { - this.res.end(); + this.response.end(); } catch { /* response may already be finished */ } @@ -355,37 +286,105 @@ class HttpSseConnection extends Connection { } } +/** + * A long-lived {@link Connection} that streams a single subscription over + * Server-Sent Events (`text/event-stream`). + * + * Unlike {@link HttpRequestConnection} (one request → one reply), this + * connection stays open: it drives the dispatcher with exactly one + * `subscribe` frame and then relays every `{event:'message'}` delivery to + * the client as an SSE `message` event until the client disconnects. + * + * The dispatcher is unchanged — it sees the same `subscribe` frame and the + * same `{event:'message', subId, payload}` deliveries it would send over + * WebSocket. This class only reframes them as SSE lines. + * + * HTTP status semantics are preserved for the *establishment* of the + * stream: the `200 text/event-stream` headers are not written until the + * dispatcher acknowledges the subscribe (`{ok:true}`). If the subscribe is + * rejected (e.g. `not_exposed`), the headers have not been sent yet, so the + * failure surfaces as a normal JSON error response with the mapped status + * code — exactly like `call`/`publish`. + */ +class HttpSubscriptionConnection extends HttpSseConnection { + /** + * @param {import('http').IncomingMessage} request + * @param {import('http').ServerResponse} response + * @param {string} capability ROS name to subscribe to. + * @param {object} [options] + * @param {number} [options.keepAliveMs=15000] Heartbeat comment interval. + */ + constructor(request, response, capability, options = {}) { + super(request, response, options); + this._capability = capability; + // Fixed subId: an SSE connection carries exactly one subscription. + this._subId = 'sse'; + request.on('close', () => this._emitCloseOnce()); + } + + /** + * Kick off dispatch: emit a single `subscribe` frame. The dispatcher + * replies synchronously with `{ok:true}` (→ start the stream) or an + * error (→ JSON failure response). + */ + begin() { + this.emit('message', { + id: this._subId, + kind: 'subscribe', + capability: this._capability, + }); + } + + send(frame) { + if (this._closed) return; + + // Subscription delivery — the steady-state case. + if (frame.event === 'message') { + this._ensureStream(); + this._writeEvent('message', frame.payload); + return; + } + + // Subscribe acknowledgement: switch the response into an event stream. + if (frame.ok === true) { + this._ensureStream(); + this._writeEvent('ready', { + capability: this._capability, + subId: this._subId, + }); + return; + } + + // Failure path. + const code = frame.code || 'internal_error'; + if (!this._streaming) { + // Headers not sent yet — surface as a normal HTTP error. + const status = _STATUS_BY_CODE[code] || 500; + this._writeError(status, code, frame.error || 'subscribe failed'); + } else { + // Already streaming — emit a terminal SSE error event, then close. + this._writeEvent('error', { error: frame.error, code }); + } + this._emitCloseOnce(); + } +} + /** Streams one action over SSE; disconnecting does not cancel the goal. */ -class HttpActionConnection extends Connection { +class HttpActionConnection extends HttpSseConnection { /** - * @param {import('http').IncomingMessage} req - * @param {import('http').ServerResponse} res + * @param {import('http').IncomingMessage} request + * @param {import('http').ServerResponse} response * @param {string} capability ROS action name. * @param {*} payload Parsed goal. * @param {object} [options] * @param {number} [options.keepAliveMs=15000] Heartbeat interval. */ - constructor(req, res, capability, payload, options = {}) { - super(); - this.req = req; - this.res = res; + constructor(request, response, capability, payload, options = {}) { + super(request, response, options); this._capability = capability; this._payload = payload; - this._keepAliveMs = - options.keepAliveMs != null ? options.keepAliveMs : 15000; - this._keepAlive = null; - this._streaming = false; - this._closed = false; // Each HTTP connection owns one goal, so a fixed ID is sufficient. this._goalId = 'http-action'; - - // Use res: req.close marks POST body completion, not client disconnect. - const onClose = () => this._emitCloseOnce(); - res.on('close', onClose); - res.on('error', (e) => { - debug('action sse response error: %s', e.message); - this._emitCloseOnce(); - }); } begin() { @@ -412,7 +411,10 @@ class HttpActionConnection extends Connection { if (frame.ok === false) { this._writeEvent('error', { error: frame.error, code: frame.code }); } else { - this._writeEvent('result', frame.payload, { status: frame.status }); + this._writeEvent('result', { + status: frame.status, + payload: frame.payload, + }); } return this._emitCloseOnce(); } @@ -429,85 +431,6 @@ class HttpActionConnection extends Connection { this._writeError(status, code, frame.error || 'send_goal failed'); this._emitCloseOnce(); } - - // eslint-disable-next-line no-unused-vars - close(code, reason) { - this._emitCloseOnce(); - } - - _ensureStream() { - if (this._streaming) return; - this._streaming = true; - this.res.writeHead(200, { - 'content-type': 'text/event-stream; charset=utf-8', - 'cache-control': 'no-cache, no-transform', - connection: 'keep-alive', - 'x-accel-buffering': 'no', - }); - if (typeof this.res.flushHeaders === 'function') { - this.res.flushHeaders(); - } - if (this._keepAliveMs > 0) { - this._keepAlive = setInterval(() => { - if (this._closed) return; - try { - this.res.write(': keep-alive\n\n'); - } catch (e) { - debug('action sse keep-alive write failed: %s', e.message); - this._emitCloseOnce(); - } - }, this._keepAliveMs); - if (typeof this._keepAlive.unref === 'function') { - this._keepAlive.unref(); - } - } - } - - _writeEvent(event, data, extra) { - if (this._closed) return; - const json = JSON.stringify( - extra ? { ...extra, payload: data } : (data ?? null) - ); - let chunk = `event: ${event}\n`; - for (const line of json.split('\n')) { - chunk += `data: ${line}\n`; - } - chunk += '\n'; - try { - this.res.write(chunk); - } catch (e) { - debug('action sse write failed: %s', e.message); - this._emitCloseOnce(); - } - } - - _writeError(status, code, message) { - const body = JSON.stringify({ ok: false, error: message, code }); - try { - this.res.writeHead(status, { - 'content-type': 'application/json; charset=utf-8', - 'content-length': Buffer.byteLength(body), - }); - this.res.end(body); - } catch (e) { - debug('action sse writeError failed: %s', e.message); - } - } - - _emitCloseOnce() { - if (this._closed) return; - this._closed = true; - if (this._keepAlive) { - clearInterval(this._keepAlive); - this._keepAlive = null; - } - try { - this.res.end(); - } catch { - /* response may already be finished */ - } - this.emit('close'); - } } /** @@ -566,7 +489,7 @@ class HttpTransport extends TransportAdapter { * request's `Origin` is echoed back when it matches). Required for a * browser on a different origin to `fetch()` / `EventSource()` this * transport. - * @param {(req: import('http').IncomingMessage) => boolean} [options.verifyRequest] + * @param {(request: import('http').IncomingMessage) => boolean} [options.verifyRequest] * Optional auth hook called with the raw request. Return `false` to * reject the request with 401. Mirrors `WebSocketTransport.verifyClient`. */ @@ -592,7 +515,9 @@ class HttpTransport extends TransportAdapter { } this._onConnection = onConnection; return new Promise((resolve, reject) => { - const server = http.createServer((req, res) => this._route(req, res)); + const server = http.createServer((request, response) => + this._route(request, response) + ); this._server = server; server.on('error', reject); server.listen(this.port, this.host, () => { @@ -626,17 +551,17 @@ class HttpTransport extends TransportAdapter { // ---------- internals ---------- - _route(req, res) { + _route(request, response) { // Apply CORS headers (if configured) before anything else, so they // appear on every response — JSON replies, 204s, SSE streams, and // errors alike. setHeader persists through later writeHead() calls. - this._applyCors(req, res); + this._applyCors(request, response); let pathname; try { - pathname = new URL(req.url || '/', 'http://localhost').pathname; + pathname = new URL(request.url || '/', 'http://localhost').pathname; } catch { - return _writeJson(res, 400, { + return _writeJson(response, 400, { ok: false, error: 'invalid request URL', code: 'invalid_url', @@ -644,7 +569,7 @@ class HttpTransport extends TransportAdapter { } if (!pathname.startsWith(this.basePath + '/')) { - return _writeJson(res, 404, { + return _writeJson(response, 404, { ok: false, error: `not a capability route: ${pathname}`, code: 'not_found', @@ -654,25 +579,25 @@ class HttpTransport extends TransportAdapter { // Preflight: browsers send OPTIONS before a cross-origin POST with a // JSON content-type. Answer it directly (no auth, no body) — but only // for capability routes, so unrelated paths still 404 above. - if (req.method === 'OPTIONS' && this.cors) { - res.writeHead(204).end(); + if (request.method === 'OPTIONS' && this.cors) { + response.writeHead(204).end(); return; } if (this.verifyRequest) { let allowed; try { - allowed = this.verifyRequest(req); + allowed = this.verifyRequest(request); } catch (e) { debug('verifyRequest threw: %s', (e && e.stack) || e); - return _writeJson(res, 500, { + return _writeJson(response, 500, { ok: false, error: 'verifyRequest hook failed', code: 'internal_error', }); } if (allowed === false) { - return _writeJson(res, 401, { + return _writeJson(response, 401, { ok: false, error: 'unauthorized', code: 'unauthorized', @@ -685,7 +610,7 @@ class HttpTransport extends TransportAdapter { // ROS name (which itself can contain slashes). const slash = tail.indexOf('/'); if (slash <= 0) { - return _writeJson(res, 404, { + return _writeJson(response, 404, { ok: false, error: `expected ${this.basePath}//`, code: 'not_found', @@ -696,7 +621,7 @@ class HttpTransport extends TransportAdapter { try { name = '/' + decodeURIComponent(tail.slice(slash + 1)); } catch { - return _writeJson(res, 400, { + return _writeJson(response, 400, { ok: false, error: `invalid percent-encoding in capability name: ${tail.slice(slash + 1)}`, code: 'invalid_url', @@ -705,22 +630,22 @@ class HttpTransport extends TransportAdapter { if (kind === 'subscribe') { if (!this.sse) { - return _writeJson(res, 404, { + return _writeJson(response, 404, { ok: false, error: 'subscribe over HTTP is disabled (enable `sse` or use WebSocket)', code: 'unsupported_kind', }); } - if (req.method !== 'GET') { - res.setHeader('allow', 'GET'); - return _writeJson(res, 405, { + if (request.method !== 'GET') { + response.setHeader('allow', 'GET'); + return _writeJson(response, 405, { ok: false, - error: `method not allowed: ${req.method} (use GET for SSE subscribe)`, + error: `method not allowed: ${request.method} (use GET for SSE subscribe)`, code: 'method_not_allowed', }); } - const conn = new HttpSseConnection(req, res, name, { + const conn = new HttpSubscriptionConnection(request, response, name, { keepAliveMs: this.sseKeepAliveMs, }); try { @@ -734,34 +659,38 @@ class HttpTransport extends TransportAdapter { } if (kind !== 'call' && kind !== 'publish' && kind !== 'action') { - return _writeJson(res, 404, { + return _writeJson(response, 404, { ok: false, error: `unsupported kind over HTTP: ${kind}`, code: 'unsupported_kind', }); } - if (req.method !== 'POST') { - res.setHeader('allow', 'POST'); - return _writeJson(res, 405, { + if (request.method !== 'POST') { + response.setHeader('allow', 'POST'); + return _writeJson(response, 405, { ok: false, - error: `method not allowed: ${req.method}`, + error: `method not allowed: ${request.method}`, code: 'method_not_allowed', }); } - _readJsonBody(req, _MAX_BODY_BYTES, (err, payload) => { + _readJsonBody(request, _MAX_BODY_BYTES, (err, payload) => { if (err) { const code = err.code || 'invalid_json'; const status = code === 'payload_too_large' ? 413 : 400; - return _writeJson(res, status, { ok: false, error: err.message, code }); + return _writeJson(response, status, { + ok: false, + error: err.message, + code, + }); } const conn = kind === 'action' - ? new HttpActionConnection(req, res, name, payload, { + ? new HttpActionConnection(request, response, name, payload, { keepAliveMs: this.sseKeepAliveMs, }) - : new HttpRequestConnection(req, res, kind, name, payload); + : new HttpRequestConnection(request, response, kind, name, payload); try { this._onConnection(conn); conn.begin(); @@ -784,44 +713,44 @@ class HttpTransport extends TransportAdapter { * rather than being limited to `content-type`. The reflected value is * added to `Vary` to keep caches correct. */ - _applyCors(req, res) { + _applyCors(request, response) { if (!this.cors) return; const allow = this.cors; // true → any; otherwise a Set of origins if (allow === true) { - res.setHeader('access-control-allow-origin', '*'); + response.setHeader('access-control-allow-origin', '*'); } else { - const origin = req.headers.origin; - res.setHeader('vary', 'Origin'); + const origin = request.headers.origin; + response.setHeader('vary', 'Origin'); if (origin && allow.has(origin)) { - res.setHeader('access-control-allow-origin', origin); + response.setHeader('access-control-allow-origin', origin); } } - res.setHeader('access-control-allow-methods', 'GET, POST, OPTIONS'); + response.setHeader('access-control-allow-methods', 'GET, POST, OPTIONS'); // Echo whatever headers the browser says it will send; fall back to // `content-type` for non-preflight requests that carry no such hint. - const requested = req.headers['access-control-request-headers']; + const requested = request.headers['access-control-request-headers']; if (requested) { - res.appendHeader('vary', 'Access-Control-Request-Headers'); - res.setHeader('access-control-allow-headers', requested); + response.appendHeader('vary', 'Access-Control-Request-Headers'); + response.setHeader('access-control-allow-headers', requested); } else { - res.setHeader('access-control-allow-headers', 'content-type'); + response.setHeader('access-control-allow-headers', 'content-type'); } - res.setHeader('access-control-max-age', '86400'); + response.setHeader('access-control-max-age', '86400'); } } -function _writeJson(res, status, body) { +function _writeJson(response, status, body) { const json = JSON.stringify(body); - res.writeHead(status, { + response.writeHead(status, { 'content-type': 'application/json; charset=utf-8', 'content-length': Buffer.byteLength(json), }); - res.end(json); + response.end(json); } -function _readJsonBody(req, maxBytes, cb) { +function _readJsonBody(request, maxBytes, cb) { // Guard against double-callback: the body may exceed maxBytes (we call - // cb(err) and req.destroy()), and the destroy itself can synchronously + // cb(err) and request.destroy()), and the destroy itself can synchronously // emit 'error' which would otherwise reach cb(err) a second time. let done = false; const finish = (err, value) => { @@ -830,14 +759,14 @@ function _readJsonBody(req, maxBytes, cb) { cb(err, value); }; - const ctype = (req.headers['content-type'] || '').toLowerCase(); + const ctype = (request.headers['content-type'] || '').toLowerCase(); // Only the media-type segment matters; ignore parameters like // `; charset=utf-8` and reject sneaky values such as // `text/plain;application/json` that include the right substring // but mean something else. const mediaType = ctype.split(';')[0].trim(); if ( - req.method === 'POST' && + request.method === 'POST' && ctype !== '' && mediaType !== 'application/json' ) { @@ -852,19 +781,19 @@ function _readJsonBody(req, maxBytes, cb) { } let total = 0; const chunks = []; - req.on('data', (chunk) => { + request.on('data', (chunk) => { if (done) return; total += chunk.length; if (total > maxBytes) { const e = new Error(`request body exceeds ${maxBytes} bytes`); e.code = 'payload_too_large'; finish(e); - req.destroy(); + request.destroy(); return; } chunks.push(chunk); }); - req.on('end', () => { + request.on('end', () => { if (done) return; const raw = Buffer.concat(chunks).toString('utf8'); if (!raw) return finish(null, {}); @@ -876,7 +805,7 @@ function _readJsonBody(req, maxBytes, cb) { finish(err); } }); - req.on('error', (e) => { + request.on('error', (e) => { e.code = e.code || 'request_error'; finish(e); }); diff --git a/test/test-web-action.js b/test/test-web-action.js index dbf3603db..b8ad5f33a 100644 --- a/test/test-web-action.js +++ b/test/test-web-action.js @@ -825,101 +825,188 @@ describe('Action capability dispatch', function () { }); }); -describe('HTTP action heartbeats', function () { - let clock; - let connection; - let response; +for (const kind of ['action', 'subscription']) { + describe(`HTTP ${kind} SSE lifecycle`, function () { + const isAction = kind === 'action'; + let clock; + let connection; + let response; + + beforeEach(function () { + clock = sinon.useFakeTimers({ toFake: ['setInterval', 'clearInterval'] }); + }); - beforeEach(function () { - clock = sinon.useFakeTimers({ toFake: ['setInterval', 'clearInterval'] }); - }); + afterEach(function () { + connection?.close(); + connection = null; + clock.restore(); + }); - afterEach(function () { - connection?.close(); - connection = null; - clock.restore(); - }); + function openStream(options) { + const transport = new HttpTransport({ sse: !isAction, ...options }); + const request = Object.assign(new EventEmitter(), { + method: isAction ? 'POST' : 'GET', + url: isAction + ? '/capability/action/fibonacci' + : '/capability/subscribe/chatter', + headers: { 'content-type': 'application/json' }, + }); + response = new EventEmitter(); + response.writeHead = sinon.stub().returns(response); + response.flushHeaders = sinon.spy(); + response.write = sinon.stub().returns(true); + response.end = sinon.spy(); + transport._onConnection = (value) => { + connection = value; + }; + transport._route(request, response); + if (isAction) { + request.emit('data', Buffer.from('{"order":5}')); + request.emit('end'); + } + assert.ok(connection); + return request; + } - function openAction(options) { - const transport = new HttpTransport(options); - const request = Object.assign(new EventEmitter(), { - method: 'POST', - url: '/capability/action/fibonacci', - headers: { 'content-type': 'application/json' }, - }); - response = new EventEmitter(); - response.writeHead = sinon.stub().returns(response); - response.flushHeaders = sinon.spy(); - response.write = sinon.stub().returns(true); - response.end = sinon.spy(); - transport._onConnection = (value) => { - connection = value; - }; - transport._route(request, response); - request.emit('data', Buffer.from('{"order":5}')); - request.emit('end'); - assert.ok(connection); - } + for (const [label, options, interval] of [ + ['default', {}, 15000], + ['configured', { sseKeepAliveMs: 10 }, 10], + ]) { + it(`uses the ${label} heartbeat interval for idle streams`, function () { + openStream(options); + assert.strictEqual(clock.countTimers(), 0); + assert.ok(response.writeHead.notCalled); + connection.send({ ok: true }); + clock.tick(interval - 1); + assert.ok(!response.write.calledWith(': keep-alive\n\n')); + clock.tick(1); + assert.ok(response.write.calledWith(': keep-alive\n\n')); + assert.strictEqual(clock.countTimers(), 1); + assert.strictEqual(connection._keepAlive.hasRef(), false); + }); + } - for (const [label, options, interval] of [ - ['default', {}, 15000], - ['configured', { sseKeepAliveMs: 10 }, 10], - ]) { - it(`uses the ${label} interval for actions without feedback`, function () { - openAction(options); - assert.strictEqual(clock.countTimers(), 0); + it('disables heartbeats when sseKeepAliveMs is zero', function () { + openStream({ sseKeepAliveMs: 0 }); connection.send({ ok: true }); - clock.tick(interval - 1); + clock.tick(30000); + assert.strictEqual(clock.countTimers(), 0); assert.ok(!response.write.calledWith(': keep-alive\n\n')); - clock.tick(1); - assert.ok(response.write.calledWith(': keep-alive\n\n')); - assert.strictEqual(clock.countTimers(), 1); - assert.strictEqual(connection._keepAlive.hasRef(), false); }); - } - - it('disables heartbeats when sseKeepAliveMs is zero', function () { - openAction({ sseKeepAliveMs: 0 }); - connection.send({ ok: true }); - clock.tick(30000); - assert.strictEqual(clock.countTimers(), 0); - assert.ok(!response.write.calledWith(': keep-alive\n\n')); - }); - for (const ending of ['result', 'error', 'disconnect', 'write failure']) { - it(`clears the timer after ${ending}`, function () { - openAction({ sseKeepAliveMs: 10 }); + it('preserves acknowledgement and data event framing', function () { + openStream({ sseKeepAliveMs: 0 }); connection.send({ ok: true }); - clock.tick(10); - assert.strictEqual(clock.countTimers(), 1); - if (ending === 'disconnect') { - response.emit('close'); - } else if (ending === 'write failure') { - response.write.throws(new Error('connection lost')); + const acknowledgement = isAction + ? 'event: accepted\ndata: {"capability":"/fibonacci"}\n\n' + : 'event: ready\ndata: {"capability":"/chatter","subId":"sse"}\n\n'; + assert.ok(response.write.calledWith(acknowledgement)); + const event = isAction ? 'feedback' : 'message'; + connection.send({ event, payload: { data: 'sample' } }); + assert.ok( + response.write.calledWith( + `event: ${event}\ndata: {"data":"sample"}\n\n` + ) + ); + assert.ok(response.writeHead.calledOnce); + assert.ok(response.flushHeaders.calledOnce); + assert.ok(response.end.notCalled); + }); + + for (const ending of [ + isAction ? 'result' : 'request close', + 'error', + 'disconnect', + 'response error', + 'write failure', + 'explicit close', + ]) { + it(`clears the timer and closes once after ${ending}`, function () { + const request = openStream({ sseKeepAliveMs: 10 }); + const onClose = sinon.spy(); + connection.on('close', onClose); + connection.send({ ok: true }); clock.tick(10); - } else { + assert.strictEqual(clock.countTimers(), 1); + if (ending === 'request close') { + request.emit('close'); + } else if (ending === 'disconnect') { + response.emit('close'); + } else if (ending === 'response error') { + response.emit('error', new Error('connection lost')); + } else if (ending === 'write failure') { + response.write.throws(new Error('connection lost')); + clock.tick(10); + } else if (ending === 'explicit close') { + connection.close(); + } else { + connection.send({ + event: isAction ? 'result' : undefined, + ok: ending === 'result', + status: 'succeeded', + payload: { sequence: [] }, + code: isAction ? 'action_failed' : 'internal_error', + error: 'stream failed', + }); + } + assert.strictEqual(clock.countTimers(), 0); + assert.strictEqual(connection._keepAlive, null); + const writes = response.write.callCount; + connection.close(); + response.emit('close'); connection.send({ - event: 'result', - ok: ending === 'result', - status: 'succeeded', - payload: { sequence: [] }, - code: 'action_failed', - error: 'action failed', + event: isAction ? 'feedback' : 'message', + payload: {}, }); - } - assert.strictEqual(clock.countTimers(), 0); - const writes = response.write.callCount; + clock.tick(30); + assert.strictEqual(response.write.callCount, writes); + assert.ok(response.end.calledOnce); + assert.ok(onClose.calledOnce); + }); + } + + it('returns a JSON error without starting a stream when rejected', function () { + openStream({ sseKeepAliveMs: 10 }); + const code = isAction ? 'goal_rejected' : 'not_exposed'; + connection.send({ ok: false, code, error: 'rejected' }); clock.tick(30); - assert.strictEqual(response.write.callCount, writes); + assert.strictEqual(clock.countTimers(), 0); + assert.ok(response.writeHead.calledWith(isAction ? 409 : 404)); + assert.ok(response.write.notCalled); + assert.ok(response.flushHeaders.notCalled); + assert.deepStrictEqual(JSON.parse(response.end.firstCall.args[0]), { + ok: false, + error: 'rejected', + code, + }); }); - } - it('does not start a timer for a rejected goal', function () { - openAction({ sseKeepAliveMs: 10 }); - connection.send({ ok: false, code: 'goal_rejected', error: 'rejected' }); - clock.tick(30); - assert.strictEqual(clock.countTimers(), 0); - assert.ok(response.writeHead.calledWith(409)); - assert.ok(response.write.notCalled); + if (isAction) { + it('keeps streaming after the POST request closes and preserves the result envelope', function () { + const request = openStream({ sseKeepAliveMs: 10 }); + const onClose = sinon.spy(); + connection.on('close', onClose); + request.emit('close'); + connection.send({ ok: true }); + request.emit('close'); + clock.tick(10); + assert.ok(response.end.notCalled); + assert.ok(onClose.notCalled); + assert.ok(response.write.calledWith(': keep-alive\n\n')); + connection.send({ + event: 'result', + status: 'succeeded', + payload: { sequence: [0, 1] }, + }); + assert.ok( + response.write.calledWith( + 'event: result\ndata: {"status":"succeeded","payload":{"sequence":[0,1]}}\n\n' + ) + ); + assert.ok(response.end.calledOnce); + assert.ok(onClose.calledOnce); + assert.strictEqual(clock.countTimers(), 0); + }); + } }); -}); +}