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 }} 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/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/lib/runtime/transports/http.js b/lib/runtime/transports/http.js index 4ef7c63de..fdbf6fcb0 100644 --- a/lib/runtime/transports/http.js +++ b/lib/runtime/transports/http.js @@ -30,6 +30,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 @@ -49,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; @@ -69,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 @@ -120,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(); } @@ -162,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) { @@ -278,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(); @@ -315,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(); @@ -325,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); } @@ -343,7 +278,7 @@ class HttpSseConnection extends Connection { this._keepAlive = null; } try { - this.res.end(); + this.response.end(); } catch { /* response may already be finished */ } @@ -351,6 +286,153 @@ 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 HttpSseConnection { + /** + * @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(request, response, capability, payload, options = {}) { + super(request, response, options); + this._capability = capability; + this._payload = payload; + // Each HTTP connection owns one goal, so a fixed ID is sufficient. + this._goalId = 'http-action'; + } + + begin() { + this.emit('message', { + id: this._goalId, + kind: 'action', + op: 'send_goal', + capability: this._capability, + payload: this._payload, + }); + } + + send(frame) { + if (this._closed) return; + + if (frame.event === 'feedback') { + this._ensureStream(); + this._writeEvent('feedback', frame.payload); + return; + } + + if (frame.event === 'result') { + this._ensureStream(); + if (frame.ok === false) { + this._writeEvent('error', { error: frame.error, code: frame.code }); + } else { + this._writeEvent('result', { + status: frame.status, + payload: frame.payload, + }); + } + return this._emitCloseOnce(); + } + + if (frame.ok === true) { + this._ensureStream(); + this._writeEvent('accepted', { capability: this._capability }); + return; + } + + // 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(); + } +} + /** * HTTP adapter for the Web Runtime. * @@ -359,6 +441,13 @@ class HttpSseConnection extends Connection { * POST /capability/call/ * POST /capability/publish/ * + * Action goals use a JSON POST with an SSE response: + * + * POST /capability/action/ (text/event-stream response) + * + * 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: * @@ -400,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`. */ @@ -426,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, () => { @@ -460,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', @@ -478,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', @@ -488,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', @@ -519,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', @@ -530,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', @@ -539,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 { @@ -567,30 +658,39 @@ class HttpTransport extends TransportAdapter { return; } - if (kind !== 'call' && kind !== 'publish') { - return _writeJson(res, 404, { + if (kind !== 'call' && kind !== 'publish' && kind !== 'action') { + return _writeJson(response, 404, { ok: false, - error: `unsupported kind over HTTP: ${kind} (use WebSocket for action)`, + 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 = new HttpRequestConnection(req, res, kind, name, payload); + const conn = + kind === 'action' + ? new HttpActionConnection(request, response, name, payload, { + keepAliveMs: this.sseKeepAliveMs, + }) + : new HttpRequestConnection(request, response, kind, name, payload); try { this._onConnection(conn); conn.begin(); @@ -613,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) => { @@ -659,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' ) { @@ -681,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, {}); @@ -705,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-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 530ce2e2d..b8ad5f33a 100644 --- a/test/test-web-action.js +++ b/test/test-web-action.js @@ -6,14 +6,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. +// 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 { 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 +36,7 @@ describe('Action capability dispatch', function () { let runtime; let server; let wsUrl; + let httpUrl; function waitOpen(ws) { return new Promise((resolve, reject) => { @@ -92,11 +98,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 () { @@ -399,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: {} } }, @@ -492,5 +517,496 @@ 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('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 { + 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(); + } + }); + + 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' }); + 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())) + ); + } + }); + + 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())) + ); + } + }); + } + + 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; + 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(); + } + }); }); }); + +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'] }); + }); + + 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; + } + + 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); + }); + } + + it('disables heartbeats when sseKeepAliveMs is zero', function () { + openStream({ sseKeepAliveMs: 0 }); + connection.send({ ok: true }); + clock.tick(30000); + assert.strictEqual(clock.countTimers(), 0); + assert.ok(!response.write.calledWith(': keep-alive\n\n')); + }); + + it('preserves acknowledgement and data event framing', function () { + openStream({ sseKeepAliveMs: 0 }); + connection.send({ ok: true }); + 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); + 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: isAction ? 'feedback' : 'message', + payload: {}, + }); + 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(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, + }); + }); + + 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); + }); + } + }); +} diff --git a/web/client.js b/web/client.js index e48b1f2b8..f38807dc8 100644 --- a/web/client.js +++ b/web/client.js @@ -24,9 +24,9 @@ // 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, ws } → explicit endpoint pair. +// - http:// / https:// → HTTP for call/publish; SSE for action. +// Subscribe lazily uses a sibling WebSocket. +// - { http, ws } → HTTP for call/publish; WS for subscribe/action. let WS = globalThis.WebSocket; let _wsResolved = !!WS; @@ -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'), { @@ -456,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) { @@ -473,6 +467,7 @@ class _HttpLink { this.baseUrl = trimmed.endsWith('/capability') ? trimmed : trimmed + '/capability'; + this._actionControllers = new Set(); } async connect() { @@ -482,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) { @@ -493,6 +491,76 @@ class _HttpLink { return this._fetch('publish', capability, payload, /* expectBody */ false); } + /** 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', + }); + } + + if (!res.ok || !res.body) { + let err = {}; + try { + 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, + }); + } + + let resolveResult, rejectResult; + const result = new Promise((res2, rej2) => { + resolveResult = res2; + rejectResult = rej2; + }); + let status; + _pumpActionStream( + res.body, + onFeedback, + resolveResult, + rejectResult, + (value) => { + status = value; + }, + signal + ).finally(() => this._actionControllers.delete(controller)); + return { + goalId: _genId(), + 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 +606,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 +621,117 @@ function _connectionLostError(reconnecting = true) { ); } +/** Runs independently; terminal SSE events settle the result promise. */ +async function _pumpActionStream( + body, + onFeedback, + resolveResult, + rejectResult, + setStatus, + signal +) { + let reader; + let buffer = ''; + let skipLeadingLF = false; + try { + reader = body.getReader(); + const decoder = new TextDecoder(); + for (;;) { + signal.throwIfAborted(); + const { done, value } = await reader.read(); + signal.throwIfAborted(); + if (done) break; + 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 = buffer.indexOf('\n\n')) !== -1) { + signal.throwIfAborted(); + const chunk = buffer.slice(0, boundary); + buffer = buffer.slice(boundary + 2); + const { event, data } = _parseSseChunk(chunk); + if (event === 'feedback' && onFeedback) { + try { + onFeedback(data); + } catch (_) { + // Callback errors must not interrupt result delivery. + } + } else if (event === 'result') { + 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 + ); + return; + } else if (event === 'error') { + 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( + signal.aborted + ? signal.reason + : Object.assign(new Error(`action stream read failed: ${e.message}`), { + code: 'network_error', + }) + ); + } finally { + if (reader) { + try { + await reader.cancel(); + } catch { + /* stream may already be closed */ + } + } + if (reader) { + try { + reader.releaseLock(); + } catch { + /* stream may already be released */ + } + } + } +} + +function _parseSseChunk(chunk) { + let event = 'message'; + const dataLines = []; + 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(/^ /, '')); + } + 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). @@ -564,10 +749,11 @@ 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. - * - object `{http, ws}` → both URLs spelled out explicitly. + * - `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}` → 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 @@ -582,7 +768,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) @@ -616,7 +802,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) @@ -661,8 +847,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; @@ -747,11 +932,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. - * 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] @@ -764,6 +947,10 @@ export class RosClient { 'action(capability, payload, options): onFeedback must be a function' ); } + if (this._closed) throw new Error('connection closed'); + if (this._http && !this._wsEager) { + return this._http.action(capability, payload, { onFeedback }); + } const ws = await this._ensureWs(); return ws.action(capability, payload, { onFeedback }); } @@ -800,7 +987,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 3997d0776..0d5055585 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. + * 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. */ @@ -190,15 +191,16 @@ 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. - * - {@link ConnectEndpoints} — both URLs spelled out. + * - `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} — 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