diff --git a/.changeset/hot-transport-ws.md b/.changeset/hot-transport-ws.md new file mode 100644 index 000000000..9b7d908cc --- /dev/null +++ b/.changeset/hot-transport-ws.md @@ -0,0 +1,5 @@ +--- +"webpack-dev-middleware": minor +--- + +Choose how hot module replacement events reach the clients with `hot.transport`: Server-Sent Events (the default), a WebSocket, or a transport of your own diff --git a/README.md b/README.md index 9ebbab998..a8d5eb281 100644 --- a/README.md +++ b/README.md @@ -326,26 +326,107 @@ middleware(compiler, { hot: true }); The object form accepts these options: -| Name | Type | Default | Description | -| :------------------------------------: | :-------: | :----------------: | :---------------------------------------------------- | -| **[`path`](#hotpath)** | `string` | `'/__webpack_hmr'` | Path the SSE endpoint is served at. | -| **[`heartbeat`](#hotheartbeat)** | `number` | `10000` | Interval (in milliseconds) between keep-alive frames. | -| **[`progress`](#hotprogress)** | `boolean` | `false` | Publish compilation progress events to the clients. | -| **[`statsOptions`](#hotstatsoptions)** | `object` | `undefined` | Deprecated — do not use; see [`stats`](#stats). | +| Name | Type | Default | Description | +| :------------------------------------: | :------------------: | :----------------: | :---------------------------------------------------- | +| **[`transport`](#hottransport)** | `string \| function` | `'sse'` | How events reach the clients. | +| **[`path`](#hotpath)** | `string` | `'/__webpack_hmr'` | Path the endpoint is served at. | +| **[`heartbeat`](#hotheartbeat)** | `number` | `10000` | Interval (in milliseconds) between keep-alive frames. | +| **[`server`](#hotserver)** | `object` | `undefined` | HTTP server the `'ws'` transport answers upgrades on. | +| **[`progress`](#hotprogress)** | `boolean` | `false` | Publish compilation progress events to the clients. | +| **[`statsOptions`](#hotstatsoptions)** | `object` | `undefined` | Deprecated — do not use; see [`stats`](#stats). | + +#### `hot.transport` + +Type: `'sse' | 'ws' | Function` +Default: `'sse'` + +How events reach the clients. + +`'sse'` serves them as [Server-Sent Events](https://developer.mozilla.org/en-US/docs/Web/API/Server-sent_events) from the middleware itself, which needs nothing else. + +`'ws'` serves them over a WebSocket. It needs the optional [`ws`](https://www.npmjs.com/package/ws) package (`npm install ws`), and an HTTP server to answer upgrades on — a handshake is an upgrade the server answers, which the middleware never sees. Give it [`hot.server`](#hotserver), or hand the server over later with the middleware's [`attach`](#attach) method: + +```js +const server = http.createServer(instance); + +instance.attach(server); +``` + +A plain `GET` on the path under `'ws'` answers `426 Upgrade Required`. + +A **function** builds a transport of your own. It is called with the resolved `path` and `heartbeat` and a logger, and must return a client stream — the same calls the built-in two answer: + +```js +/** + * @param {{ path: string, heartbeat: number }} options + * @param {Logger} logger + * @returns {ClientStream} + */ +middleware(compiler, { + hot: { + transport: ({ path, heartbeat }, logger) => ({ + // Answer a request on the endpoint's path. + handler(req, res) {}, + // True while at least one client is connected; a compilation with no + // clients skips serializing its payload. + hasClients: () => clients.size > 0, + // Call `fn` with each client once it has joined. It is what catches a + // client up with the last hashes, so it can apply the next update. + onConnect(fn) {}, + // Publish a payload to every client. + publish(payload) {}, + // Publish a payload to one client, as handed to `onConnect`. + publishTo(client, payload) {}, + // End every client and stop any timers. + close() {}, + // Optional, for transports built on an upgrade. + attach(server) {}, + detach() {}, + }), + }, +}); +``` + +A function that returns something missing one of those throws, naming what is absent, rather than failing later from wherever it is first published to. + +The clients are yours — whatever `onConnect` hands out is what `publishTo` takes back — so in TypeScript name their type through `ClientStreamFactory`: + +```ts +import { type ClientStreamFactory } from "webpack-dev-middleware/types/hot"; + +interface MyClient { + id: number; + send: (frame: string) => void; +} + +const transport: ClientStreamFactory = ({ path }, logger) => ({ + // ... + publishTo(client, payload) { + client.send(JSON.stringify(payload)); + }, +}); +``` #### `hot.path` Type: `String` Default: `'/__webpack_hmr'` -Path the SSE endpoint is served at. Must start with a slash and match the `path` option used by the client. +Path the endpoint is served at. Must start with a slash and match the `path` option used by the client. #### `hot.heartbeat` Type: `Number` Default: `10000` -Heartbeat interval (in milliseconds) used to keep the SSE connection alive when no compilation events are produced. Must be `1` or greater. +Heartbeat interval (in milliseconds) used to keep the connection alive when no compilation events are produced: keep-alive frames for `'sse'`, pings for `'ws'`. Must be `1` or greater. + +#### `hot.server` + +Type: `Object` +Default: `undefined` + +HTTP server the [`'ws'`](#hottransport) transport answers upgrades on, when it already exists where the middleware is built. Otherwise hand it over later with the middleware's [`attach`](#attach) method. Ignored by `'sse'`, which is answered by the middleware itself. #### `hot.progress` @@ -582,6 +663,42 @@ hotClient.subscribe((payload) => { `webpack-dev-middleware` also provides convenience methods that can be use to interact with the middleware at runtime: +### `attach(server)` + +Gives the [`hot.transport: "ws"`](#hottransport) endpoint the HTTP server to answer WebSocket upgrades on. A handshake is an upgrade the server answers, which the middleware never sees, so it cannot find the server on its own. Use this when the server is built after the middleware; when it already exists, [`hot.server`](#hotserver) does the same thing. + +Does nothing when `hot` is disabled or the transport is Server-Sent Events, which the middleware answers itself. + +#### Parameters + +##### `server` + +Type: `http.Server | https.Server` +Required: `Yes` + +The server whose `upgrade` event the endpoint listens on. It stops listening when the middleware is closed. + +```js +const http = require("node:http"); +const express = require("express"); +const webpack = require("webpack"); + +const middleware = require("webpack-dev-middleware"); + +const compiler = webpack({/* Webpack configuration */}); +const instance = middleware(compiler, { hot: { transport: "ws" } }); + +// eslint-disable-next-line new-cap +const app = new express(); + +app.use(instance); + +const server = http.createServer(app); + +instance.attach(server); +server.listen(3000); +``` + ### `close(callback)` Instructs `webpack-dev-middleware` instance to stop watching for file changes. diff --git a/package-lock.json b/package-lock.json index 4ed146090..150b53a9a 100644 --- a/package-lock.json +++ b/package-lock.json @@ -28,6 +28,7 @@ "@types/express": "^5.0.6", "@types/mime-types": "^3.0.1", "@types/node": "^26.4.1", + "@types/ws": "^8.18.1", "acorn": "^8.18.0", "ajv": "^8.20.0", "babel-jest": "^30.1.2", @@ -59,7 +60,8 @@ "router": "^2.2.0", "supertest": "^7.2.2", "typescript": "^6.0.3", - "webpack": "^5.110.3" + "webpack": "^5.110.3", + "ws": "^8.18.3" }, "engines": { "node": ">= 20.9.0" @@ -68,6 +70,9 @@ "type": "opencollective", "url": "https://opencollective.com/webpack" }, + "optionalDependencies": { + "ws": "^8.18.3" + }, "peerDependencies": { "webpack": "^5.101.0" }, @@ -5819,6 +5824,16 @@ "dev": true, "license": "MIT" }, + "node_modules/@types/ws": { + "version": "8.18.1", + "resolved": "https://registry.npmjs.org/@types/ws/-/ws-8.18.1.tgz", + "integrity": "sha512-ThVF6DCVhA8kUGy+aazFQ4kXQ7E1Ty7A3ypFOe0IcJV8O/M511G99AW24irKrW56Wt44yG9+ij8FaqoBGkuBXg==", + "dev": true, + "license": "MIT", + "dependencies": { + "@types/node": "*" + } + }, "node_modules/@types/yargs": { "version": "17.0.35", "resolved": "https://registry.npmjs.org/@types/yargs/-/yargs-17.0.35.tgz", @@ -6244,9 +6259,6 @@ "arm64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -6261,9 +6273,6 @@ "arm64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -6278,9 +6287,6 @@ "loong64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -6295,9 +6301,6 @@ "loong64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -6312,9 +6315,6 @@ "ppc64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -6329,9 +6329,6 @@ "riscv64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -6346,9 +6343,6 @@ "riscv64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -6363,9 +6357,6 @@ "s390x" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -6380,9 +6371,6 @@ "x64" ], "dev": true, - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -6397,9 +6385,6 @@ "x64" ], "dev": true, - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ diff --git a/package.json b/package.json index 16b39a889..488969c72 100644 --- a/package.json +++ b/package.json @@ -82,6 +82,7 @@ "@types/express": "^5.0.6", "@types/mime-types": "^3.0.1", "@types/node": "^26.4.1", + "@types/ws": "^8.18.1", "acorn": "^8.18.0", "ajv": "^8.20.0", "babel-jest": "^30.1.2", @@ -113,7 +114,8 @@ "router": "^2.2.0", "supertest": "^7.2.2", "typescript": "^6.0.3", - "webpack": "^5.110.3" + "webpack": "^5.110.3", + "ws": "^8.18.3" }, "peerDependencies": { "webpack": "^5.101.0" @@ -123,6 +125,9 @@ "optional": true } }, + "optionalDependencies": { + "ws": "^8.18.3" + }, "engines": { "node": ">= 20.9.0" } diff --git a/src/hot.js b/src/hot.js index bc0e1b047..1ef73202b 100644 --- a/src/hot.js +++ b/src/hot.js @@ -7,6 +7,7 @@ /** @typedef {import("webpack").StatsError} StatsError */ /** @typedef {import("./index.js").IncomingMessage} IncomingMessage */ /** @typedef {import("./index.js").ServerResponse} ServerResponse */ +/** @typedef {import("node:http").Server} HttpServer */ // The object form only (no presets/booleans) — it is merged over the // middleware's own base options, which string or boolean forms cannot be. @@ -15,8 +16,10 @@ /** * @typedef {object} HotOptions - * @property {string=} path the path the SSE endpoint is served at + * @property {("sse" | "ws" | ClientStreamFactory)=} transport how events reach the clients, Server-Sent Events by default + * @property {string=} path the path the endpoint is served at * @property {number=} heartbeat heartbeat interval in milliseconds + * @property {HttpServer=} server HTTP server the `"ws"` transport answers upgrades on, when it is already built * @property {StatsOptions=} statsOptions deprecated, removed in the next major release — webpack stats options used when serializing compilation results * @property {boolean=} progress publish compilation progress events to the clients */ @@ -34,17 +37,60 @@ * @property {string[]=} errors errors */ +// eslint-disable-next-line jsdoc/reject-any-type +/** @typedef {any} EXPECTED_ANY */ + +/** + * The WebSocket members a client is published to through. Structural rather than + * the ws package's own declarations, which would put an optional dependency's + * types in the path of every consumer, including those on Server-Sent Events. + * @typedef {object} WebSocketLikeClient + * @property {number} readyState the socket's current state + * @property {number} OPEN the value `readyState` has while the socket is open + * @property {(data: string) => void} send send a frame to this client + */ + /** - * @typedef {object} EventStream - * @property {(req: IncomingMessage, res: ServerResponse) => void} handler attach a new client + * What a client is addressed by, which is whatever the transport handed out: the + * response holding a Server-Sent Events stream, or a WebSocket. + * @typedef {ServerResponse | WebSocketLikeClient} StreamClient + */ + +/** + * One transport's clients. `createHot` publishes through this and does not know + * whether the events leave over Server-Sent Events, a WebSocket or something of + * your own, which is what `TClient` is for: a transport built by a `transport` + * function names the type of the clients it hands to `onConnect` and takes back + * in `publishTo`. + * @template {EXPECTED_ANY} [TClient=StreamClient] + * @typedef {object} ClientStream + * @property {(req: IncomingMessage, res: ServerResponse) => void} handler answer a request on the endpoint's path * @property {() => boolean} hasClients true when at least one client is connected + * @property {(fn: (client: TClient) => void) => void} onConnect called with each client once it has joined * @property {(payload: Payload | { action: string }) => void} publish publish a payload to every client - * @property {(res: ServerResponse, payload: Payload | { action: string }) => void} publishTo publish a payload to a single client + * @property {(client: TClient, payload: Payload | { action: string }) => void} publishTo publish a payload to a single client * @property {() => void} close end every client and stop the heartbeat + * @property {((server: HttpServer) => void)=} attach answer upgrades on this server + * @property {(() => void)=} detach stop answering upgrades + */ + +/** + * Builds a transport of your own. The same calls `createHot` makes of the + * built-in two are made of whatever this returns. + * @template {EXPECTED_ANY} [TClient=StreamClient] + * @callback ClientStreamFactory + * @param {{ path: string, heartbeat: number }} options the endpoint's path and heartbeat interval + * @param {Logger} logger logger + * @returns {ClientStream} client stream */ +/** @typedef {ClientStream} EventStream */ + +const createWebSocketStream = require("./servers/WebSocketServer.js"); + const HOT_DEFAULT_PATH = "/__webpack_hmr"; const HOT_DEFAULT_HEARTBEAT = 10 * 1000; +const HOT_DEFAULT_TRANSPORT = "sse"; const PLUGIN_NAME = "DevMiddleware"; /** @@ -62,6 +108,41 @@ function pathMatch(url, expected) { } } +// Everything `createHot` calls on a stream. A transport of your own which is +// missing one of them would throw from wherever it is first published to, +// which is a long way from the option that built it. +const CLIENT_STREAM_METHODS = [ + "close", + "handler", + "hasClients", + "onConnect", + "publish", + "publishTo", +]; + +/** + * @param {ClientStream} stream what a `transport` function returned + * @returns {ClientStream} the same stream + */ +function checkClientStream(stream) { + const missing = + stream && typeof stream === "object" + ? CLIENT_STREAM_METHODS.filter( + (method) => + typeof (/** @type {Record} */ (stream)[method]) !== + "function", + ) + : CLIENT_STREAM_METHODS; + + if (missing.length > 0) { + throw new TypeError( + `The 'hot.transport' function must return a client stream, which is missing: ${missing.join(", ")}.`, + ); + } + + return stream; +} + /** * @param {number} heartbeat heartbeat interval in milliseconds * @param {Logger} logger logger @@ -71,6 +152,8 @@ function createEventStream(heartbeat, logger) { let clientId = 0; /** @type {Map} */ let clients = new Map(); + /** @type {((client: StreamClient) => void) | undefined} */ + let onConnectFn; /** * Run the callback for every client that can still be written to — a @@ -125,6 +208,9 @@ function createEventStream(heartbeat, logger) { hasClients() { return clients.size > 0; }, + onConnect(fn) { + onConnectFn = fn; + }, handler(req, res) { // A response another middleware already started can no longer become an // SSE stream — end it instead of crashing on writeHead. @@ -187,6 +273,12 @@ function createEventStream(heartbeat, logger) { // again, so it would stay in `clients` forever. if (req.destroyed) { disconnect(); + + return; + } + + if (onConnectFn) { + onConnectFn(res); } }, publish(payload) { @@ -201,7 +293,9 @@ function createEventStream(heartbeat, logger) { client.write(frame); }); }, - publishTo(res, payload) { + publishTo(client, payload) { + const res = /** @type {ServerResponse} */ (client); + if (res.writableEnded) { return; } @@ -383,8 +477,10 @@ function publishBundles(bundles, previousBundles, eventStream) { /** * @typedef {object} HotInstance - * @property {string} path path the SSE endpoint is served at - * @property {(req: IncomingMessage, res: ServerResponse) => void} handle attach the request as a SSE client + * @property {string} path path the endpoint is served at + * @property {("sse" | "ws" | ClientStreamFactory)} transport how events reach the clients + * @property {(server: HttpServer) => void} attach answer WebSocket upgrades on this server, a no-op for Server-Sent Events + * @property {(req: IncomingMessage, res: ServerResponse) => void} handle answer a request on the endpoint's path * @property {(payload: Payload | { action: string }) => void} publish publish a payload to every client * @property {() => void} close end every client and detach the heartbeat */ @@ -399,6 +495,7 @@ function createHot(compiler, userOptions, statsOption) { const options = userOptions === true ? {} : userOptions; const path = options.path || HOT_DEFAULT_PATH; const heartbeat = options.heartbeat ?? HOT_DEFAULT_HEARTBEAT; + const transport = options.transport || HOT_DEFAULT_TRANSPORT; const { statsOptions } = options; const logger = compiler.getInfrastructureLogger("webpack-dev-middleware"); @@ -409,8 +506,25 @@ function createHot(compiler, userOptions, statsOption) { ); } - let eventStream = createEventStream(heartbeat, logger); - logger.log(`Hot module replacement enabled, serving events at "${path}"`); + /** @type {ClientStream} */ + let eventStream; + /** @type {string} */ + let transportName; + + if (typeof transport === "function") { + eventStream = checkClientStream(transport({ heartbeat, path }, logger)); + transportName = "a custom transport"; + } else if (transport === "ws") { + eventStream = createWebSocketStream({ heartbeat, path }, logger); + transportName = "a WebSocket"; + } else { + eventStream = createEventStream(heartbeat, logger); + transportName = "Server-Sent Events"; + } + + logger.log( + `Hot module replacement enabled, serving events at "${path}" over ${transportName}`, + ); // `latestBundles` survives rebuilds so hashes can be compared per build. /** @type {StatsCompilation[] | null} */ @@ -419,6 +533,24 @@ function createHot(compiler, userOptions, statsOption) { let closed = false; let lastProgressPercent = -1; + // Catch a new client up wherever it joined from, as `sync` events carrying + // the last hashes. + eventStream.onConnect((client) => { + if (!valid || !latestBundles) { + return; + } + + for (const stats of latestBundles) { + eventStream.publishTo(client, bundlePayload(stats, "sync")); + } + }); + + // A WebSocket is upgraded by the HTTP server rather than answered by the + // middleware, so the transport needs the server itself. + if (options.server && eventStream.attach) { + eventStream.attach(options.server); + } + if (options.progress) { const { webpack } = "compilers" in compiler ? compiler.compilers[0] : compiler; @@ -494,6 +626,14 @@ function createHot(compiler, userOptions, statsOption) { return { path, + transport, + attach(server) { + if (closed || !eventStream.attach) { + return; + } + + eventStream.attach(server); + }, handle(req, res) { // A request can race `close()` past the middleware intercept — end it // instead of leaving it hanging without a response. @@ -505,13 +645,6 @@ function createHot(compiler, userOptions, statsOption) { } eventStream.handler(req, res); - - // Catch only the new client up, as `sync` events with the last hashes. - if (valid && latestBundles) { - for (const stats of latestBundles) { - eventStream.publishTo(res, bundlePayload(stats, "sync")); - } - } }, publish(payload) { if (closed) return; @@ -525,7 +658,9 @@ function createHot(compiler, userOptions, statsOption) { // https://github.com/webpack/tapable/issues/32#issuecomment-350644466 closed = true; eventStream.close(); - eventStream = /** @type {EventStream} */ (/** @type {unknown} */ (null)); + eventStream = /** @type {ClientStream} */ ( + /** @type {unknown} */ (null) + ); }, }; } @@ -533,6 +668,8 @@ function createHot(compiler, userOptions, statsOption) { module.exports = createHot; module.exports.HOT_DEFAULT_HEARTBEAT = HOT_DEFAULT_HEARTBEAT; module.exports.HOT_DEFAULT_PATH = HOT_DEFAULT_PATH; +module.exports.HOT_DEFAULT_TRANSPORT = HOT_DEFAULT_TRANSPORT; +module.exports.checkClientStream = checkClientStream; module.exports.createEventStream = createEventStream; module.exports.createHot = createHot; module.exports.formatErrors = formatErrors; diff --git a/src/index.js b/src/index.js index eab161152..840b0ba57 100644 --- a/src/index.js +++ b/src/index.js @@ -150,6 +150,11 @@ const noop = () => {}; * @param {Callback} callback */ +/** + * @callback Attach + * @param {import("node:http").Server} server HTTP server the `hot.transport: "ws"` endpoint answers upgrades on + */ + /** * @callback Close * @param {(err: Error | null | undefined) => void} callback @@ -162,6 +167,7 @@ const noop = () => {}; * @property {GetFilenameFromUrl} getFilenameFromUrl get filename from url * @property {WaitUntilValid} waitUntilValid wait until valid * @property {Invalidate} invalidate invalidate + * @property {Attach} attach answer WebSocket upgrades on this server * @property {Close} close close * @property {Context} context context */ @@ -625,6 +631,14 @@ function wdm(compiler, options = {}, isPlugin = false) { instance.getFilenameFromUrl = (url) => middleware.getFilenameFromUrl(filledContext, url); + // A WebSocket handshake is an upgrade the HTTP server answers, which the + // middleware never sees, so the server is handed over rather than inferred. + instance.attach = (server) => { + if (filledContext.hot) { + filledContext.hot.attach(server); + } + }; + instance.waitUntilValid = (callback = noop) => { middleware.ready(filledContext, callback); }; diff --git a/src/options.check.js b/src/options.check.js index c1cd0c6d1..64d9b13ec 100644 --- a/src/options.check.js +++ b/src/options.check.js @@ -2,4 +2,4 @@ // DO NOT MODIFY BY HAND. Run `npm run fix:schema-check` to update. /* eslint-disable */ // @ts-nocheck -"use strict";module.exports = validate10;module.exports.default = validate10;const schema11 = {"type":"object","properties":{"mimeTypes":{"description":"Allows a user to register custom mime types or extension mappings.","link":"https://github.com/webpack/webpack-dev-middleware#mimetypes","type":"object"},"mimeTypeDefault":{"description":"Allows a user to register a default mime type when we can't determine the content type.","link":"https://github.com/webpack/webpack-dev-middleware#mimetypedefault","type":"string"},"writeToDisk":{"description":"Allows to write generated files on disk.","link":"https://github.com/webpack/webpack-dev-middleware#writetodisk","anyOf":[{"type":"boolean"},{"instanceof":"Function"}]},"methods":{"description":"Allows to pass the list of HTTP request methods accepted by the middleware.","link":"https://github.com/webpack/webpack-dev-middleware#methods","type":"array","items":{"type":"string","minLength":1}},"headers":{"anyOf":[{"type":"array","items":{"type":"object","additionalProperties":false,"properties":{"key":{"description":"key of header.","type":"string"},"value":{"description":"value of header.","type":"string"}}},"minItems":1},{"type":"object"},{"instanceof":"Function"}],"description":"Allows to pass custom HTTP headers on each request","link":"https://github.com/webpack/webpack-dev-middleware#headers"},"publicPath":{"description":"The `publicPath` specifies the public URL address of the output files when referenced in a browser.","link":"https://github.com/webpack/webpack-dev-middleware#publicpath","anyOf":[{"enum":["auto"]},{"type":"string"},{"instanceof":"Function"}]},"stats":{"description":"Stats options object or preset name.","link":"https://github.com/webpack/webpack-dev-middleware#stats","anyOf":[{"enum":["none","summary","errors-only","errors-warnings","minimal","normal","detailed","verbose"]},{"type":"boolean"},{"type":"object","additionalProperties":true}]},"serverSideRender":{"description":"Instructs the module to enable or disable the server-side rendering mode.","link":"https://github.com/webpack/webpack-dev-middleware#serversiderender","type":"boolean"},"outputFileSystem":{"description":"Set the default file system which will be used by webpack as primary destination of generated files.","link":"https://github.com/webpack/webpack-dev-middleware#outputfilesystem","type":"object"},"index":{"description":"Allows to serve an index of the directory.","link":"https://github.com/webpack/webpack-dev-middleware#index","anyOf":[{"type":"boolean"},{"type":"string","minLength":1}]},"modifyResponseData":{"description":"Allows to set up a callback to change the response data.","link":"https://github.com/webpack/webpack-dev-middleware#modifyresponsedata","instanceof":"Function"},"etag":{"description":"Enable or disable etag generation.","link":"https://github.com/webpack/webpack-dev-middleware#etag","enum":["weak","strong"]},"lastModified":{"description":"Enable or disable `Last-Modified` header. Uses the file system's last modified value.","link":"https://github.com/webpack/webpack-dev-middleware#lastmodified","type":"boolean"},"cacheControl":{"description":"Enable or disable setting `Cache-Control` response header.","link":"https://github.com/webpack/webpack-dev-middleware#cachecontrol","anyOf":[{"type":"boolean"},{"type":"number"},{"type":"string","minLength":1},{"type":"object","properties":{"maxAge":{"type":"number"},"immutable":{"type":"boolean"}},"additionalProperties":false}]},"cacheImmutable":{"description":"Enable or disable setting `Cache-Control: public, max-age=31536000, immutable` response header for immutable assets (i.e. asset with a hash in file name like `image.a4c12bde.jpg`).","link":"https://github.com/webpack/webpack-dev-middleware#cacheimmutable","type":"boolean"},"forwardError":{"description":"Enable or disable forwarding errors to next middleware.","link":"https://github.com/webpack/webpack-dev-middleware#forwarderrors","type":"boolean"},"hot":{"description":"Enable hot module replacement via a Server-Sent Events endpoint.","link":"https://github.com/webpack/webpack-dev-middleware#hot","anyOf":[{"type":"boolean"},{"type":"object","additionalProperties":false,"properties":{"path":{"description":"The path the SSE endpoint is served at. Must start with a slash and carry no query string or fragment.","type":"string","pattern":"^/[^?#]*$"},"heartbeat":{"description":"Heartbeat interval (in milliseconds) used to keep the SSE connection alive.","type":"number","minimum":1},"progress":{"description":"Publish compilation progress events to the clients.","type":"boolean"},"statsOptions":{"description":"Deprecated, do not use, will be removed in the next major release. Use the `stats` option instead, which decides whether a payload carries errors and warnings.","type":"object","additionalProperties":true}}}]}},"additionalProperties":false};const func2 = Object.prototype.hasOwnProperty;const pattern0 = new RegExp("^/[^?#]*$", "u");function validate10(data, {instancePath="", parentData, parentDataProperty, rootData=data}={}){let vErrors = null;let errors = 0;if(errors === 0){if(data && typeof data == "object" && !Array.isArray(data)){const _errs1 = errors;for(const key0 in data){if(!(func2.call(schema11.properties, key0))){validate10.errors = [{instancePath,schemaPath:"#/additionalProperties",keyword:"additionalProperties",params:{additionalProperty: key0},message:"must NOT have additional properties"}];return false;break;}}if(_errs1 === errors){if(data.mimeTypes !== undefined){let data0 = data.mimeTypes;const _errs2 = errors;if(!(data0 && typeof data0 == "object" && !Array.isArray(data0))){validate10.errors = [{instancePath:instancePath+"/mimeTypes",schemaPath:"#/properties/mimeTypes/type",keyword:"type",params:{type: "object"},message:"must be object"}];return false;}var valid0 = _errs2 === errors;}else {var valid0 = true;}if(valid0){if(data.mimeTypeDefault !== undefined){const _errs5 = errors;if(typeof data.mimeTypeDefault !== "string"){validate10.errors = [{instancePath:instancePath+"/mimeTypeDefault",schemaPath:"#/properties/mimeTypeDefault/type",keyword:"type",params:{type: "string"},message:"must be string"}];return false;}var valid0 = _errs5 === errors;}else {var valid0 = true;}if(valid0){if(data.writeToDisk !== undefined){let data2 = data.writeToDisk;const _errs8 = errors;const _errs9 = errors;let valid1 = false;const _errs10 = errors;if(typeof data2 !== "boolean"){const err0 = {instancePath:instancePath+"/writeToDisk",schemaPath:"#/properties/writeToDisk/anyOf/0/type",keyword:"type",params:{type: "boolean"},message:"must be boolean"};if(vErrors === null){vErrors = [err0];}else {vErrors.push(err0);}errors++;}var _valid0 = _errs10 === errors;valid1 = valid1 || _valid0;if(!valid1){const _errs12 = errors;if(!(data2 instanceof Function)){const err1 = {instancePath:instancePath+"/writeToDisk",schemaPath:"#/properties/writeToDisk/anyOf/1/instanceof",keyword:"instanceof",params:{},message:"must pass \"instanceof\" keyword validation"};if(vErrors === null){vErrors = [err1];}else {vErrors.push(err1);}errors++;}var _valid0 = _errs12 === errors;valid1 = valid1 || _valid0;}if(!valid1){const err2 = {instancePath:instancePath+"/writeToDisk",schemaPath:"#/properties/writeToDisk/anyOf",keyword:"anyOf",params:{},message:"must match a schema in anyOf"};if(vErrors === null){vErrors = [err2];}else {vErrors.push(err2);}errors++;validate10.errors = vErrors;return false;}else {errors = _errs9;if(vErrors !== null){if(_errs9){vErrors.length = _errs9;}else {vErrors = null;}}}var valid0 = _errs8 === errors;}else {var valid0 = true;}if(valid0){if(data.methods !== undefined){let data3 = data.methods;const _errs14 = errors;if(errors === _errs14){if(Array.isArray(data3)){var valid2 = true;const len0 = data3.length;for(let i0=0; i0=", limit: 1},message:"must be >= 1"};if(vErrors === null){vErrors = [err37];}else {vErrors.push(err37);}errors++;}}else {const err38 = {instancePath:instancePath+"/hot/heartbeat",schemaPath:"#/properties/hot/anyOf/1/properties/heartbeat/type",keyword:"type",params:{type: "number"},message:"must be number"};if(vErrors === null){vErrors = [err38];}else {vErrors.push(err38);}errors++;}}var valid12 = _errs101 === errors;}else {var valid12 = true;}if(valid12){if(data22.progress !== undefined){const _errs103 = errors;if(typeof data22.progress !== "boolean"){const err39 = {instancePath:instancePath+"/hot/progress",schemaPath:"#/properties/hot/anyOf/1/properties/progress/type",keyword:"type",params:{type: "boolean"},message:"must be boolean"};if(vErrors === null){vErrors = [err39];}else {vErrors.push(err39);}errors++;}var valid12 = _errs103 === errors;}else {var valid12 = true;}if(valid12){if(data22.statsOptions !== undefined){let data26 = data22.statsOptions;const _errs105 = errors;if(errors === _errs105){if(data26 && typeof data26 == "object" && !Array.isArray(data26)){}else {const err40 = {instancePath:instancePath+"/hot/statsOptions",schemaPath:"#/properties/hot/anyOf/1/properties/statsOptions/type",keyword:"type",params:{type: "object"},message:"must be object"};if(vErrors === null){vErrors = [err40];}else {vErrors.push(err40);}errors++;}}var valid12 = _errs105 === errors;}else {var valid12 = true;}}}}}}else {const err41 = {instancePath:instancePath+"/hot",schemaPath:"#/properties/hot/anyOf/1/type",keyword:"type",params:{type: "object"},message:"must be object"};if(vErrors === null){vErrors = [err41];}else {vErrors.push(err41);}errors++;}}var _valid6 = _errs96 === errors;valid11 = valid11 || _valid6;}if(!valid11){const err42 = {instancePath:instancePath+"/hot",schemaPath:"#/properties/hot/anyOf",keyword:"anyOf",params:{},message:"must match a schema in anyOf"};if(vErrors === null){vErrors = [err42];}else {vErrors.push(err42);}errors++;validate10.errors = vErrors;return false;}else {errors = _errs93;if(vErrors !== null){if(_errs93){vErrors.length = _errs93;}else {vErrors = null;}}}var valid0 = _errs92 === errors;}else {var valid0 = true;}}}}}}}}}}}}}}}}}}}else {validate10.errors = [{instancePath,schemaPath:"#/type",keyword:"type",params:{type: "object"},message:"must be object"}];return false;}}validate10.errors = vErrors;return errors === 0;} \ No newline at end of file +"use strict";module.exports = validate10;module.exports.default = validate10;const schema11 = {"type":"object","properties":{"mimeTypes":{"description":"Allows a user to register custom mime types or extension mappings.","link":"https://github.com/webpack/webpack-dev-middleware#mimetypes","type":"object"},"mimeTypeDefault":{"description":"Allows a user to register a default mime type when we can't determine the content type.","link":"https://github.com/webpack/webpack-dev-middleware#mimetypedefault","type":"string"},"writeToDisk":{"description":"Allows to write generated files on disk.","link":"https://github.com/webpack/webpack-dev-middleware#writetodisk","anyOf":[{"type":"boolean"},{"instanceof":"Function"}]},"methods":{"description":"Allows to pass the list of HTTP request methods accepted by the middleware.","link":"https://github.com/webpack/webpack-dev-middleware#methods","type":"array","items":{"type":"string","minLength":1}},"headers":{"anyOf":[{"type":"array","items":{"type":"object","additionalProperties":false,"properties":{"key":{"description":"key of header.","type":"string"},"value":{"description":"value of header.","type":"string"}}},"minItems":1},{"type":"object"},{"instanceof":"Function"}],"description":"Allows to pass custom HTTP headers on each request","link":"https://github.com/webpack/webpack-dev-middleware#headers"},"publicPath":{"description":"The `publicPath` specifies the public URL address of the output files when referenced in a browser.","link":"https://github.com/webpack/webpack-dev-middleware#publicpath","anyOf":[{"enum":["auto"]},{"type":"string"},{"instanceof":"Function"}]},"stats":{"description":"Stats options object or preset name.","link":"https://github.com/webpack/webpack-dev-middleware#stats","anyOf":[{"enum":["none","summary","errors-only","errors-warnings","minimal","normal","detailed","verbose"]},{"type":"boolean"},{"type":"object","additionalProperties":true}]},"serverSideRender":{"description":"Instructs the module to enable or disable the server-side rendering mode.","link":"https://github.com/webpack/webpack-dev-middleware#serversiderender","type":"boolean"},"outputFileSystem":{"description":"Set the default file system which will be used by webpack as primary destination of generated files.","link":"https://github.com/webpack/webpack-dev-middleware#outputfilesystem","type":"object"},"index":{"description":"Allows to serve an index of the directory.","link":"https://github.com/webpack/webpack-dev-middleware#index","anyOf":[{"type":"boolean"},{"type":"string","minLength":1}]},"modifyResponseData":{"description":"Allows to set up a callback to change the response data.","link":"https://github.com/webpack/webpack-dev-middleware#modifyresponsedata","instanceof":"Function"},"etag":{"description":"Enable or disable etag generation.","link":"https://github.com/webpack/webpack-dev-middleware#etag","enum":["weak","strong"]},"lastModified":{"description":"Enable or disable `Last-Modified` header. Uses the file system's last modified value.","link":"https://github.com/webpack/webpack-dev-middleware#lastmodified","type":"boolean"},"cacheControl":{"description":"Enable or disable setting `Cache-Control` response header.","link":"https://github.com/webpack/webpack-dev-middleware#cachecontrol","anyOf":[{"type":"boolean"},{"type":"number"},{"type":"string","minLength":1},{"type":"object","properties":{"maxAge":{"type":"number"},"immutable":{"type":"boolean"}},"additionalProperties":false}]},"cacheImmutable":{"description":"Enable or disable setting `Cache-Control: public, max-age=31536000, immutable` response header for immutable assets (i.e. asset with a hash in file name like `image.a4c12bde.jpg`).","link":"https://github.com/webpack/webpack-dev-middleware#cacheimmutable","type":"boolean"},"forwardError":{"description":"Enable or disable forwarding errors to next middleware.","link":"https://github.com/webpack/webpack-dev-middleware#forwarderrors","type":"boolean"},"hot":{"description":"Enable hot module replacement over a Server-Sent Events or WebSocket endpoint.","link":"https://github.com/webpack/webpack-dev-middleware#hot","anyOf":[{"type":"boolean"},{"type":"object","additionalProperties":false,"properties":{"transport":{"description":"How events reach the clients: `sse`, `ws` (needs the optional `ws` dependency and an HTTP server to answer upgrades on, given as `server` or through the middleware's `attach` method), or a function building a transport of your own.","anyOf":[{"enum":["sse","ws"]},{"instanceof":"Function"}]},"path":{"description":"The path the endpoint is served at. Must start with a slash and carry no query string or fragment.","type":"string","pattern":"^/[^?#]*$"},"heartbeat":{"description":"Heartbeat interval (in milliseconds) used to keep the connection alive.","type":"number","minimum":1},"server":{"description":"HTTP server the `ws` transport answers upgrades on, when it is already built. Otherwise hand it over later with the middleware's `attach` method.","type":"object","additionalProperties":true},"progress":{"description":"Publish compilation progress events to the clients.","type":"boolean"},"statsOptions":{"description":"Deprecated, do not use, will be removed in the next major release. Use the `stats` option instead, which decides whether a payload carries errors and warnings.","type":"object","additionalProperties":true}}}]}},"additionalProperties":false};const func2 = Object.prototype.hasOwnProperty;const pattern0 = new RegExp("^/[^?#]*$", "u");function validate10(data, {instancePath="", parentData, parentDataProperty, rootData=data}={}){let vErrors = null;let errors = 0;if(errors === 0){if(data && typeof data == "object" && !Array.isArray(data)){const _errs1 = errors;for(const key0 in data){if(!(func2.call(schema11.properties, key0))){validate10.errors = [{instancePath,schemaPath:"#/additionalProperties",keyword:"additionalProperties",params:{additionalProperty: key0},message:"must NOT have additional properties"}];return false;break;}}if(_errs1 === errors){if(data.mimeTypes !== undefined){let data0 = data.mimeTypes;const _errs2 = errors;if(!(data0 && typeof data0 == "object" && !Array.isArray(data0))){validate10.errors = [{instancePath:instancePath+"/mimeTypes",schemaPath:"#/properties/mimeTypes/type",keyword:"type",params:{type: "object"},message:"must be object"}];return false;}var valid0 = _errs2 === errors;}else {var valid0 = true;}if(valid0){if(data.mimeTypeDefault !== undefined){const _errs5 = errors;if(typeof data.mimeTypeDefault !== "string"){validate10.errors = [{instancePath:instancePath+"/mimeTypeDefault",schemaPath:"#/properties/mimeTypeDefault/type",keyword:"type",params:{type: "string"},message:"must be string"}];return false;}var valid0 = _errs5 === errors;}else {var valid0 = true;}if(valid0){if(data.writeToDisk !== undefined){let data2 = data.writeToDisk;const _errs8 = errors;const _errs9 = errors;let valid1 = false;const _errs10 = errors;if(typeof data2 !== "boolean"){const err0 = {instancePath:instancePath+"/writeToDisk",schemaPath:"#/properties/writeToDisk/anyOf/0/type",keyword:"type",params:{type: "boolean"},message:"must be boolean"};if(vErrors === null){vErrors = [err0];}else {vErrors.push(err0);}errors++;}var _valid0 = _errs10 === errors;valid1 = valid1 || _valid0;if(!valid1){const _errs12 = errors;if(!(data2 instanceof Function)){const err1 = {instancePath:instancePath+"/writeToDisk",schemaPath:"#/properties/writeToDisk/anyOf/1/instanceof",keyword:"instanceof",params:{},message:"must pass \"instanceof\" keyword validation"};if(vErrors === null){vErrors = [err1];}else {vErrors.push(err1);}errors++;}var _valid0 = _errs12 === errors;valid1 = valid1 || _valid0;}if(!valid1){const err2 = {instancePath:instancePath+"/writeToDisk",schemaPath:"#/properties/writeToDisk/anyOf",keyword:"anyOf",params:{},message:"must match a schema in anyOf"};if(vErrors === null){vErrors = [err2];}else {vErrors.push(err2);}errors++;validate10.errors = vErrors;return false;}else {errors = _errs9;if(vErrors !== null){if(_errs9){vErrors.length = _errs9;}else {vErrors = null;}}}var valid0 = _errs8 === errors;}else {var valid0 = true;}if(valid0){if(data.methods !== undefined){let data3 = data.methods;const _errs14 = errors;if(errors === _errs14){if(Array.isArray(data3)){var valid2 = true;const len0 = data3.length;for(let i0=0; i0=", limit: 1},message:"must be >= 1"};if(vErrors === null){vErrors = [err40];}else {vErrors.push(err40);}errors++;}}else {const err41 = {instancePath:instancePath+"/hot/heartbeat",schemaPath:"#/properties/hot/anyOf/1/properties/heartbeat/type",keyword:"type",params:{type: "number"},message:"must be number"};if(vErrors === null){vErrors = [err41];}else {vErrors.push(err41);}errors++;}}var valid12 = _errs105 === errors;}else {var valid12 = true;}if(valid12){if(data22.server !== undefined){let data26 = data22.server;const _errs107 = errors;if(errors === _errs107){if(data26 && typeof data26 == "object" && !Array.isArray(data26)){}else {const err42 = {instancePath:instancePath+"/hot/server",schemaPath:"#/properties/hot/anyOf/1/properties/server/type",keyword:"type",params:{type: "object"},message:"must be object"};if(vErrors === null){vErrors = [err42];}else {vErrors.push(err42);}errors++;}}var valid12 = _errs107 === errors;}else {var valid12 = true;}if(valid12){if(data22.progress !== undefined){const _errs110 = errors;if(typeof data22.progress !== "boolean"){const err43 = {instancePath:instancePath+"/hot/progress",schemaPath:"#/properties/hot/anyOf/1/properties/progress/type",keyword:"type",params:{type: "boolean"},message:"must be boolean"};if(vErrors === null){vErrors = [err43];}else {vErrors.push(err43);}errors++;}var valid12 = _errs110 === errors;}else {var valid12 = true;}if(valid12){if(data22.statsOptions !== undefined){let data28 = data22.statsOptions;const _errs112 = errors;if(errors === _errs112){if(data28 && typeof data28 == "object" && !Array.isArray(data28)){}else {const err44 = {instancePath:instancePath+"/hot/statsOptions",schemaPath:"#/properties/hot/anyOf/1/properties/statsOptions/type",keyword:"type",params:{type: "object"},message:"must be object"};if(vErrors === null){vErrors = [err44];}else {vErrors.push(err44);}errors++;}}var valid12 = _errs112 === errors;}else {var valid12 = true;}}}}}}}}else {const err45 = {instancePath:instancePath+"/hot",schemaPath:"#/properties/hot/anyOf/1/type",keyword:"type",params:{type: "object"},message:"must be object"};if(vErrors === null){vErrors = [err45];}else {vErrors.push(err45);}errors++;}}var _valid6 = _errs96 === errors;valid11 = valid11 || _valid6;}if(!valid11){const err46 = {instancePath:instancePath+"/hot",schemaPath:"#/properties/hot/anyOf",keyword:"anyOf",params:{},message:"must match a schema in anyOf"};if(vErrors === null){vErrors = [err46];}else {vErrors.push(err46);}errors++;validate10.errors = vErrors;return false;}else {errors = _errs93;if(vErrors !== null){if(_errs93){vErrors.length = _errs93;}else {vErrors = null;}}}var valid0 = _errs92 === errors;}else {var valid0 = true;}}}}}}}}}}}}}}}}}}}else {validate10.errors = [{instancePath,schemaPath:"#/type",keyword:"type",params:{type: "object"},message:"must be object"}];return false;}}validate10.errors = vErrors;return errors === 0;} \ No newline at end of file diff --git a/src/options.json b/src/options.json index 3a71d0397..883e748ed 100644 --- a/src/options.json +++ b/src/options.json @@ -179,7 +179,7 @@ "type": "boolean" }, "hot": { - "description": "Enable hot module replacement via a Server-Sent Events endpoint.", + "description": "Enable hot module replacement over a Server-Sent Events or WebSocket endpoint.", "link": "https://github.com/webpack/webpack-dev-middleware#hot", "anyOf": [ { @@ -189,16 +189,32 @@ "type": "object", "additionalProperties": false, "properties": { + "transport": { + "description": "How events reach the clients: `sse`, `ws` (needs the optional `ws` dependency and an HTTP server to answer upgrades on, given as `server` or through the middleware's `attach` method), or a function building a transport of your own.", + "anyOf": [ + { + "enum": ["sse", "ws"] + }, + { + "instanceof": "Function" + } + ] + }, "path": { - "description": "The path the SSE endpoint is served at. Must start with a slash and carry no query string or fragment.", + "description": "The path the endpoint is served at. Must start with a slash and carry no query string or fragment.", "type": "string", "pattern": "^/[^?#]*$" }, "heartbeat": { - "description": "Heartbeat interval (in milliseconds) used to keep the SSE connection alive.", + "description": "Heartbeat interval (in milliseconds) used to keep the connection alive.", "type": "number", "minimum": 1 }, + "server": { + "description": "HTTP server the `ws` transport answers upgrades on, when it is already built. Otherwise hand it over later with the middleware's `attach` method.", + "type": "object", + "additionalProperties": true + }, "progress": { "description": "Publish compilation progress events to the clients.", "type": "boolean" diff --git a/src/servers/WebSocketServer.js b/src/servers/WebSocketServer.js new file mode 100644 index 000000000..d51d3b9eb --- /dev/null +++ b/src/servers/WebSocketServer.js @@ -0,0 +1,229 @@ +/** @typedef {import("node:http").Server} HttpServer */ +/** @typedef {import("node:http").IncomingMessage} IncomingMessage */ +/** @typedef {import("node:stream").Duplex} Duplex */ +/** @typedef {import("ws").WebSocket} WebSocket */ +/** @typedef {typeof import("ws").WebSocketServer} WsServerConstructor */ +/** @typedef {import("../hot.js").Logger} Logger */ +/** @typedef {import("../hot.js").Payload} Payload */ +/** @typedef {import("../hot.js").ClientStream} ClientStream */ + +// How often a client is pinged to find out whether it is still there. A client +// that has not answered the previous ping is dropped rather than pinged again. +const WS_DEFAULT_HEARTBEAT = 10 * 1000; + +/** + * `ws` is only needed by `transport: "ws"`, so it is an optional dependency and + * is required here rather than at the top of the module. + * @returns {WsServerConstructor} the `ws` server constructor + */ +function requireWsServer() { + try { + return require("ws").WebSocketServer; + } catch { + throw new Error( + "The 'hot.transport: \"ws\"' option needs the 'ws' package, which is an optional dependency of webpack-dev-middleware. Install it with `npm install ws`, or use the default 'hot.transport: \"sse\"'.", + ); + } +} + +/** + * A client stream carried over WebSocket rather than Server-Sent Events. It + * answers the same calls as `createEventStream`, so `createHot` does not know + * which of them it is publishing to. + * @param {object} options options + * @param {string} options.path the path the endpoint is served at + * @param {number} options.heartbeat heartbeat interval in milliseconds + * @param {Logger} logger logger + * @returns {ClientStream} client stream + */ +function createWebSocketStream({ path, heartbeat }, logger) { + const WebSocketServerImplementation = requireWsServer(); + /** @type {Set} */ + const clients = new Set(); + /** @type {((client: WebSocket) => void) | undefined} */ + let onConnectFn; + /** @type {HttpServer | undefined} */ + let attachedServer; + /** @type {((req: IncomingMessage, socket: Duplex, head: Buffer) => void) | undefined} */ + let upgradeListener; + + const implementation = new WebSocketServerImplementation({ + noServer: true, + path, + // `clients` is tracked here so a client is dropped the moment it closes, + // which is what `hasClients` reads. + clientTracking: false, + }); + + implementation.on( + "error", + /** @param {Error} err error */ (err) => { + logger.error(err.message); + }, + ); + + // A client that did not answer the previous ping is gone: a half-open socket + // never emits `close`, so nothing else would ever remove it. + /** @type {WeakSet} */ + let awaitingPong = new WeakSet(); + /** @type {ReturnType | null} */ + let interval = null; + + const startHeartbeat = () => { + if (interval !== null) { + return; + } + + interval = setInterval(() => { + for (const client of clients) { + if (awaitingPong.has(client)) { + client.terminate(); + continue; + } + + awaitingPong.add(client); + client.ping(() => {}); + } + }, heartbeat); + + // Don't block process exit on the heartbeat timer. + if (typeof interval.unref === "function") { + interval.unref(); + } + }; + + const stopHeartbeat = () => { + if (interval !== null) { + clearInterval(interval); + interval = null; + } + }; + + implementation.on( + "connection", + /** @param {WebSocket} client client */ (client) => { + clients.add(client); + startHeartbeat(); + logger.log(`Client connected (${clients.size} active)`); + + client.on("pong", () => { + awaitingPong.delete(client); + }); + + client.on("close", () => { + clients.delete(client); + awaitingPong.delete(client); + + if (clients.size === 0) { + stopHeartbeat(); + } + + logger.log(`Client disconnected (${clients.size} active)`); + }); + + client.on( + "error", + /** @param {Error} err error */ (err) => { + logger.error(err.message); + }, + ); + + if (onConnectFn) { + onConnectFn(client); + } + }, + ); + + const detach = () => { + if (attachedServer && upgradeListener) { + attachedServer.removeListener("upgrade", upgradeListener); + } + + attachedServer = undefined; + upgradeListener = undefined; + }; + + return { + attach(server) { + // Attaching twice would upgrade every request twice over. + if (attachedServer === server) { + return; + } + + if (attachedServer) { + detach(); + } + + upgradeListener = (req, socket, head) => { + // Another WebSocket endpoint on the same server owns this path. + if (!implementation.shouldHandle(req)) { + return; + } + + implementation.handleUpgrade(req, socket, head, (client) => { + implementation.emit("connection", client, req); + }); + }; + + attachedServer = server; + server.on("upgrade", upgradeListener); + }, + close() { + stopHeartbeat(); + detach(); + + for (const client of clients) { + client.close(); + } + + clients.clear(); + awaitingPong = new WeakSet(); + implementation.close(); + }, + detach, + handler(req, res) { + // The handshake is an upgrade the HTTP server answers, so a plain request + // reaching the middleware is a client which cannot speak this transport. + if (!res.headersSent) { + res.writeHead(426, { "Content-Type": "text/plain; charset=utf-8" }); + } + + if (!res.writableEnded) { + res.end("Upgrade Required"); + } + }, + hasClients() { + return clients.size > 0; + }, + onConnect(fn) { + onConnectFn = fn; + }, + publish(payload) { + // With no clients connected there is nothing to serialize for. + if (clients.size === 0) { + return; + } + + const frame = JSON.stringify(payload); + + for (const client of clients) { + if (client.readyState === client.OPEN) { + client.send(frame); + } + } + }, + publishTo(client, payload) { + const socket = /** @type {WebSocket} */ (client); + + if (socket.readyState !== socket.OPEN) { + return; + } + + socket.send(JSON.stringify(payload)); + }, + }; +} + +module.exports = createWebSocketStream; +module.exports.WS_DEFAULT_HEARTBEAT = WS_DEFAULT_HEARTBEAT; +module.exports.createWebSocketStream = createWebSocketStream; diff --git a/test/__snapshots__/validation-options.test.js.snap.webpack5 b/test/__snapshots__/validation-options.test.js.snap.webpack5 index f37594501..28769c7b8 100644 --- a/test/__snapshots__/validation-options.test.js.snap.webpack5 +++ b/test/__snapshots__/validation-options.test.js.snap.webpack5 @@ -78,37 +78,37 @@ exports[`validation should throw an error on the "headers" option with "true" va exports[`validation should throw an error on the "hot" option with "{"heartbeat":-1}" value 1`] = ` "Invalid options object. Dev Middleware has been initialized using an options object that does not match the API schema. - options.hot.heartbeat should be >= 1. - -> Heartbeat interval (in milliseconds) used to keep the SSE connection alive." + -> Heartbeat interval (in milliseconds) used to keep the connection alive." `; exports[`validation should throw an error on the "hot" option with "{"heartbeat":0}" value 1`] = ` "Invalid options object. Dev Middleware has been initialized using an options object that does not match the API schema. - options.hot.heartbeat should be >= 1. - -> Heartbeat interval (in milliseconds) used to keep the SSE connection alive." + -> Heartbeat interval (in milliseconds) used to keep the connection alive." `; exports[`validation should throw an error on the "hot" option with "{"path":""}" value 1`] = ` "Invalid options object. Dev Middleware has been initialized using an options object that does not match the API schema. - options.hot.path should match pattern "^/[^?#]*$". - -> The path the SSE endpoint is served at. Must start with a slash and carry no query string or fragment." + -> The path the endpoint is served at. Must start with a slash and carry no query string or fragment." `; exports[`validation should throw an error on the "hot" option with "{"path":"/__hmr#section"}" value 1`] = ` "Invalid options object. Dev Middleware has been initialized using an options object that does not match the API schema. - options.hot.path should match pattern "^/[^?#]*$". - -> The path the SSE endpoint is served at. Must start with a slash and carry no query string or fragment." + -> The path the endpoint is served at. Must start with a slash and carry no query string or fragment." `; exports[`validation should throw an error on the "hot" option with "{"path":"/__hmr?client=1"}" value 1`] = ` "Invalid options object. Dev Middleware has been initialized using an options object that does not match the API schema. - options.hot.path should match pattern "^/[^?#]*$". - -> The path the SSE endpoint is served at. Must start with a slash and carry no query string or fragment." + -> The path the endpoint is served at. Must start with a slash and carry no query string or fragment." `; exports[`validation should throw an error on the "hot" option with "{"path":"hmr"}" value 1`] = ` "Invalid options object. Dev Middleware has been initialized using an options object that does not match the API schema. - options.hot.path should match pattern "^/[^?#]*$". - -> The path the SSE endpoint is served at. Must start with a slash and carry no query string or fragment." + -> The path the endpoint is served at. Must start with a slash and carry no query string or fragment." `; exports[`validation should throw an error on the "hot" option with "{"statsOptions":"errors-only"}" value 1`] = ` @@ -125,34 +125,66 @@ exports[`validation should throw an error on the "hot" option with "{"statsOptio -> Deprecated, do not use, will be removed in the next major release. Use the \`stats\` option instead, which decides whether a payload carries errors and warnings." `; +exports[`validation should throw an error on the "hot" option with "{"transport":"websocket"}" value 1`] = ` +"Invalid options object. Dev Middleware has been initialized using an options object that does not match the API schema. + - options.hot should be one of these: + boolean | object { transport?, path?: string, heartbeat?: number, server?, progress?: boolean, statsOptions? } + -> Enable hot module replacement over a Server-Sent Events or WebSocket endpoint. + -> Read more at https://github.com/webpack/webpack-dev-middleware#hot + Details: + * options.hot.transport should be one of these: + "sse" | "ws" | function + -> How events reach the clients: \`sse\`, \`ws\` (needs the optional \`ws\` dependency and an HTTP server to answer upgrades on, given as \`server\` or through the middleware's \`attach\` method), or a function building a transport of your own. + Details: + * options.hot.transport should be one of these: + "sse" | "ws" + * options.hot.transport should be an instance of function." +`; + +exports[`validation should throw an error on the "hot" option with "{"transport":true}" value 1`] = ` +"Invalid options object. Dev Middleware has been initialized using an options object that does not match the API schema. + - options.hot should be one of these: + boolean | object { transport?, path?: string, heartbeat?: number, server?, progress?: boolean, statsOptions? } + -> Enable hot module replacement over a Server-Sent Events or WebSocket endpoint. + -> Read more at https://github.com/webpack/webpack-dev-middleware#hot + Details: + * options.hot.transport should be one of these: + "sse" | "ws" | function + -> How events reach the clients: \`sse\`, \`ws\` (needs the optional \`ws\` dependency and an HTTP server to answer upgrades on, given as \`server\` or through the middleware's \`attach\` method), or a function building a transport of your own. + Details: + * options.hot.transport should be one of these: + "sse" | "ws" + * options.hot.transport should be an instance of function." +`; + exports[`validation should throw an error on the "hot" option with "{"unknown":true}" value 1`] = ` "Invalid options object. Dev Middleware has been initialized using an options object that does not match the API schema. - options.hot has an unknown property 'unknown'. These properties are valid: - object { path?: string, heartbeat?: number, progress?: boolean, statsOptions? }" + object { transport?, path?: string, heartbeat?: number, server?, progress?: boolean, statsOptions? }" `; exports[`validation should throw an error on the "hot" option with "0" value 1`] = ` "Invalid options object. Dev Middleware has been initialized using an options object that does not match the API schema. - options.hot should be one of these: - boolean | object { path?: string, heartbeat?: number, progress?: boolean, statsOptions? } - -> Enable hot module replacement via a Server-Sent Events endpoint. + boolean | object { transport?, path?: string, heartbeat?: number, server?, progress?: boolean, statsOptions? } + -> Enable hot module replacement over a Server-Sent Events or WebSocket endpoint. -> Read more at https://github.com/webpack/webpack-dev-middleware#hot Details: * options.hot should be a boolean. * options.hot should be an object: - object { path?: string, heartbeat?: number, progress?: boolean, statsOptions? }" + object { transport?, path?: string, heartbeat?: number, server?, progress?: boolean, statsOptions? }" `; exports[`validation should throw an error on the "hot" option with "foo" value 1`] = ` "Invalid options object. Dev Middleware has been initialized using an options object that does not match the API schema. - options.hot should be one of these: - boolean | object { path?: string, heartbeat?: number, progress?: boolean, statsOptions? } - -> Enable hot module replacement via a Server-Sent Events endpoint. + boolean | object { transport?, path?: string, heartbeat?: number, server?, progress?: boolean, statsOptions? } + -> Enable hot module replacement over a Server-Sent Events or WebSocket endpoint. -> Read more at https://github.com/webpack/webpack-dev-middleware#hot Details: * options.hot should be a boolean. * options.hot should be an object: - object { path?: string, heartbeat?: number, progress?: boolean, statsOptions? }" + object { transport?, path?: string, heartbeat?: number, server?, progress?: boolean, statsOptions? }" `; exports[`validation should throw an error on the "index" option with "{}" value 1`] = ` diff --git a/test/hot.test.js b/test/hot.test.js index cb8150612..0b8ff0214 100644 --- a/test/hot.test.js +++ b/test/hot.test.js @@ -502,8 +502,10 @@ describe("createHot", () => { const hot = createHot(compiler, { path: "/__hmr" }); + // Which transport answered is part of it: the two are configured the same + // way but fail in different places. expect(messages).toContain( - 'Hot module replacement enabled, serving events at "/__hmr"', + 'Hot module replacement enabled, serving events at "/__hmr" over Server-Sent Events', ); hot.close(); @@ -1030,3 +1032,297 @@ describe("createHot", () => { expect(writes).toHaveLength(0); }); }); + +describe("createHot over a WebSocket", () => { + /** @type {(() => Promise)[]} */ + let cleanups = []; + + // Registered rather than left to each test, so a failing assertion still + // tears the server down instead of leaving the run hanging on it. + afterEach(async () => { + for (const cleanup of cleanups.reverse()) { + await cleanup(); + } + + cleanups = []; + }); + + /** + * Stand up an HTTP server with a `ws` hot instance attached to it. + * @param {EXPECTED_OBJECT} compiler fake compiler + * @param {EXPECTED_OBJECT=} options extra hot options + * @returns {Promise<{ hot: EXPECTED_OBJECT, url: string, stop: () => Promise }>} the running endpoint + */ + async function serveOverWs(compiler, options = {}) { + const hot = createHot(compiler, { transport: "ws", ...options }); + const server = http.createServer((req, res) => { + hot.handle(req, res); + }); + + hot.attach(server); + + await new Promise((resolve) => { + server.listen(0, resolve); + }); + + const { port } = server.address(); + + const stop = async () => { + hot.close(); + await new Promise((resolve) => { + server.close(resolve); + }); + }; + + cleanups.push(stop); + + return { hot, url: `ws://127.0.0.1:${port}${hot.path}`, stop }; + } + + /** + * Connect a client and collect the payloads it is sent. + * @param {string} url url to connect to + * @returns {Promise<{ socket: EXPECTED_OBJECT, messages: string[] }>} the open client + */ + async function connect(url) { + const { WebSocket } = require("ws"); + + const socket = new WebSocket(url); + /** @type {string[]} */ + const messages = []; + + socket.on("message", (data) => messages.push(data.toString())); + + cleanups.push(async () => { + socket.terminate(); + }); + + await new Promise((resolve, reject) => { + // A server which never answers the upgrade leaves the socket pending + // rather than erroring, so the wait is bounded here. + const timer = setTimeout(() => { + socket.terminate(); + reject(new Error(`timed out connecting to ${url}`)); + }, 5000); + + socket.on("open", () => { + clearTimeout(timer); + resolve(); + }); + socket.on("error", (error) => { + clearTimeout(timer); + reject(error); + }); + }); + + return { socket, messages }; + } + + /** + * @param {() => boolean} predicate what to wait for + * @returns {Promise} resolves once the predicate holds + */ + async function until(predicate) { + const deadline = Date.now() + 5000; + + while (!predicate()) { + if (Date.now() > deadline) { + throw new Error("timed out"); + } + + await new Promise((resolve) => { + setTimeout(resolve, 10); + }); + } + } + + it("logs that the events are served over a WebSocket", async () => { + const messages = []; + const compiler = makeFakeCompiler({ + log: (message) => messages.push(message), + error: () => {}, + }); + + await serveOverWs(compiler, { path: "/__hmr" }); + + expect(messages).toContain( + 'Hot module replacement enabled, serving events at "/__hmr" over a WebSocket', + ); + }); + + it("publishes a build to a connected client", async () => { + const compiler = makeFakeCompiler(); + const endpoint = await serveOverWs(compiler); + const { messages } = await connect(endpoint.url); + + compiler.emitDone(makeFakeStats()); + await until(() => messages.length > 0); + + // The same payloads the SSE transport publishes, without its `data:` frame. + expect(JSON.parse(messages[0]).action).toBe("built"); + }); + + it("catches a client up with sync when it joins after a build", async () => { + const compiler = makeFakeCompiler(); + const endpoint = await serveOverWs(compiler); + + compiler.emitDone(makeFakeStats({ hash: "joined-late" })); + + const { messages } = await connect(endpoint.url); + + await until(() => messages.length > 0); + + // A client that missed the build still has to learn the current hash, or + // it can never apply the next update. + const sync = JSON.parse(messages[0]); + + expect(sync.action).toBe("sync"); + expect(sync.hash).toBe("joined-late"); + }); + + it("answers a plain request on the path with 426 rather than a stream", async () => { + const compiler = makeFakeCompiler(); + const endpoint = await serveOverWs(compiler); + const url = endpoint.url.replace("ws://", "http://"); + + const response = await fetch(url); + + // Nothing but an upgrade can speak this transport, so a client which asked + // for the path over plain HTTP is told so instead of being left hanging. + expect(response.status).toBe(426); + }); + + it("stops answering upgrades once closed", async () => { + const compiler = makeFakeCompiler(); + const endpoint = await serveOverWs(compiler); + + endpoint.hot.close(); + + await expect(connect(endpoint.url)).rejects.toThrow( + /timed out connecting|ECONNREFUSED|socket hang up|Unexpected server response/, + ); + }); +}); + +describe("createHot over a transport of your own", () => { + /** + * The smallest thing `createHot` will publish through: it keeps the payloads + * in an array rather than putting them on a wire. + * @returns {EXPECTED_OBJECT} a recording client stream + */ + function createRecordingStream() { + /** @type {EXPECTED_OBJECT[]} */ + const clients = []; + /** @type {EXPECTED_OBJECT[]} */ + const published = []; + /** @type {EXPECTED_OBJECT} */ + let onConnectFn; + + return { + clients, + published, + close() {}, + connect(client) { + clients.push(client); + onConnectFn(client); + }, + handler() {}, + hasClients: () => clients.length > 0, + onConnect(fn) { + onConnectFn = fn; + }, + publish(payload) { + published.push(payload); + }, + publishTo(client, payload) { + client.sent.push(payload); + }, + }; + } + + it("publishes through it instead of the built-in transports", () => { + const compiler = makeFakeCompiler(); + const stream = createRecordingStream(); + const hot = createHot(compiler, { transport: () => stream }); + + stream.connect({ sent: [] }); + compiler.emitDone(makeFakeStats()); + + expect(stream.published.map((p) => p.action)).toContain("built"); + + hot.close(); + }); + + it("is handed the resolved path and heartbeat", () => { + const compiler = makeFakeCompiler(); + /** @type {EXPECTED_OBJECT} */ + let received; + const hot = createHot(compiler, { + heartbeat: 1234, + path: "/__custom", + transport: (options) => { + received = options; + + return createRecordingStream(); + }, + }); + + // Resolved, not raw: a transport should not have to re-apply the defaults. + expect(received).toEqual({ heartbeat: 1234, path: "/__custom" }); + + hot.close(); + }); + + it("catches a late joiner up with sync, the same as the built-in ones", () => { + const compiler = makeFakeCompiler(); + const stream = createRecordingStream(); + const hot = createHot(compiler, { transport: () => stream }); + + compiler.emitDone(makeFakeStats({ hash: "late" })); + + const client = { sent: [] }; + + stream.connect(client); + + expect(client.sent).toHaveLength(1); + expect(client.sent[0].action).toBe("sync"); + expect(client.sent[0].hash).toBe("late"); + + hot.close(); + }); + + it("keeps the optional ws dependency out of the published types", () => { + const fs = require("node:fs"); + const path = require("node:path"); + + const declarations = fs.readFileSync( + path.resolve(__dirname, "../types/hot.d.ts"), + "utf8", + ); + + // `types/index.d.ts` reaches these, so an `import("ws")` here is loaded by + // every consumer — including one on Server-Sent Events, who has no reason + // to have `@types/ws` installed, and whose build then fails outright. + expect(declarations).not.toMatch(/\bimport\("ws"\)/); + }); + + it("names what a returned object is missing rather than failing later", () => { + const compiler = makeFakeCompiler(); + + // Publishing through a half-built stream would throw from whichever call + // it lacks, a long way from the option that built it. + expect(() => + createHot(compiler, { transport: () => ({ publish() {} }) }), + ).toThrow( + "The 'hot.transport' function must return a client stream, which is missing: close, handler, hasClients, onConnect, publishTo.", + ); + }); + + it("treats a transport returning nothing the same way", () => { + const compiler = makeFakeCompiler(); + + expect(() => createHot(compiler, { transport: () => undefined })).toThrow( + /must return a client stream, which is missing: close, handler/, + ); + }); +}); diff --git a/test/validation-options.test.js b/test/validation-options.test.js index 2341ee189..6c6411a28 100644 --- a/test/validation-options.test.js +++ b/test/validation-options.test.js @@ -92,6 +92,20 @@ describe("validation", () => { {}, { path: "/__hmr" }, { heartbeat: 1000 }, + { transport: "sse" }, + { transport: "ws" }, + // A transport of your own, built by this function. It is called, so it + // has to answer what `createHot` publishes through. + { + transport: () => ({ + close() {}, + handler() {}, + hasClients: () => false, + onConnect() {}, + publish() {}, + publishTo() {}, + }), + }, { statsOptions: { all: false } }, ], failure: [ @@ -103,6 +117,8 @@ describe("validation", () => { { path: "hmr" }, { path: "/__hmr?client=1" }, { path: "/__hmr#section" }, + { transport: "websocket" }, + { transport: true }, { heartbeat: -1 }, // 0 would silently fall back to the default interval — reject it. { heartbeat: 0 }, diff --git a/types/hot.d.ts b/types/hot.d.ts index a19f68e0b..45f563556 100644 --- a/types/hot.d.ts +++ b/types/hot.d.ts @@ -1,8 +1,10 @@ export = createHot; /** * @typedef {object} HotInstance - * @property {string} path path the SSE endpoint is served at - * @property {(req: IncomingMessage, res: ServerResponse) => void} handle attach the request as a SSE client + * @property {string} path path the endpoint is served at + * @property {("sse" | "ws" | ClientStreamFactory)} transport how events reach the clients + * @property {(server: HttpServer) => void} attach answer WebSocket upgrades on this server, a no-op for Server-Sent Events + * @property {(req: IncomingMessage, res: ServerResponse) => void} handle answer a request on the endpoint's path * @property {(payload: Payload | { action: string }) => void} publish publish a payload to every client * @property {() => void} close end every client and detach the heartbeat */ @@ -21,6 +23,8 @@ declare namespace createHot { export { HOT_DEFAULT_HEARTBEAT, HOT_DEFAULT_PATH, + HOT_DEFAULT_TRANSPORT, + checkClientStream, createEventStream, createHot, formatErrors, @@ -37,53 +41,29 @@ declare namespace createHot { StatsError, IncomingMessage, ServerResponse, + HttpServer, StatsOptions, MiddlewareStatsOption, HotOptions, Payload, + EXPECTED_ANY, + WebSocketLikeClient, + StreamClient, + ClientStream, + ClientStreamFactory, EventStream, }; } declare const HOT_DEFAULT_HEARTBEAT: number; -/** @typedef {import("webpack").Compiler} Compiler */ -/** @typedef {import("webpack").MultiCompiler} MultiCompiler */ -/** @typedef {ReturnType} Logger */ -/** @typedef {import("webpack").Stats} Stats */ -/** @typedef {import("webpack").MultiStats} MultiStats */ -/** @typedef {import("webpack").StatsCompilation} StatsCompilation */ -/** @typedef {import("webpack").StatsError} StatsError */ -/** @typedef {import("./index.js").IncomingMessage} IncomingMessage */ -/** @typedef {import("./index.js").ServerResponse} ServerResponse */ -/** @typedef {import("webpack").StatsOptions} StatsOptions */ -/** @typedef {import("webpack").Configuration["stats"]} MiddlewareStatsOption */ -/** - * @typedef {object} HotOptions - * @property {string=} path the path the SSE endpoint is served at - * @property {number=} heartbeat heartbeat interval in milliseconds - * @property {StatsOptions=} statsOptions deprecated, removed in the next major release — webpack stats options used when serializing compilation results - * @property {boolean=} progress publish compilation progress events to the clients - */ -/** - * @typedef {object} Payload - * @property {string} action action - * @property {string=} file file that invalidated the compilation - * @property {string=} name name - * @property {number=} time time - * @property {string=} hash hash - * @property {number=} percent compilation progress (0-100) - * @property {string=} message progress message - * @property {string[]=} warnings warnings - * @property {string[]=} errors errors - */ +declare const HOT_DEFAULT_PATH: "/__webpack_hmr"; +declare const HOT_DEFAULT_TRANSPORT: "sse"; /** - * @typedef {object} EventStream - * @property {(req: IncomingMessage, res: ServerResponse) => void} handler attach a new client - * @property {() => boolean} hasClients true when at least one client is connected - * @property {(payload: Payload | { action: string }) => void} publish publish a payload to every client - * @property {(res: ServerResponse, payload: Payload | { action: string }) => void} publishTo publish a payload to a single client - * @property {() => void} close end every client and stop the heartbeat + * @param {ClientStream} stream what a `transport` function returned + * @returns {ClientStream} the same stream */ -declare const HOT_DEFAULT_PATH: "/__webpack_hmr"; +declare function checkClientStream( + stream: ClientStream, +): ClientStream; /** * @param {number} heartbeat heartbeat interval in milliseconds * @param {Logger} logger logger @@ -130,11 +110,19 @@ declare function toBundles( ): StatsCompilation[]; type HotInstance = { /** - * path the SSE endpoint is served at + * path the endpoint is served at */ path: string; /** - * attach the request as a SSE client + * how events reach the clients + */ + transport: "sse" | "ws" | ClientStreamFactory; + /** + * answer WebSocket upgrades on this server, a no-op for Server-Sent Events + */ + attach: (server: HttpServer) => void; + /** + * answer a request on the endpoint's path */ handle: (req: IncomingMessage, res: ServerResponse) => void; /** @@ -161,17 +149,26 @@ type StatsCompilation = import("webpack").StatsCompilation; type StatsError = import("webpack").StatsError; type IncomingMessage = import("./index.js").IncomingMessage; type ServerResponse = import("./index.js").ServerResponse; +type HttpServer = import("node:http").Server; type StatsOptions = import("webpack").StatsOptions; type MiddlewareStatsOption = import("webpack").Configuration["stats"]; type HotOptions = { /** - * the path the SSE endpoint is served at + * how events reach the clients, Server-Sent Events by default + */ + transport?: ("sse" | "ws" | ClientStreamFactory) | undefined; + /** + * the path the endpoint is served at */ path?: string | undefined; /** * heartbeat interval in milliseconds */ heartbeat?: number | undefined; + /** + * HTTP server the `"ws"` transport answers upgrades on, when it is already built + */ + server?: HttpServer | undefined; /** * deprecated, removed in the next major release — webpack stats options used when serializing compilation results */ @@ -219,15 +216,51 @@ type Payload = { */ errors?: string[] | undefined; }; -type EventStream = { +type EXPECTED_ANY = any; +/** + * The WebSocket members a client is published to through. Structural rather than + * the ws package's own declarations, which would put an optional dependency's + * types in the path of every consumer, including those on Server-Sent Events. + */ +type WebSocketLikeClient = { + /** + * the socket's current state + */ + readyState: number; /** - * attach a new client + * the value `readyState` has while the socket is open + */ + OPEN: number; + /** + * send a frame to this client + */ + send: (data: string) => void; +}; +/** + * What a client is addressed by, which is whatever the transport handed out: the + * response holding a Server-Sent Events stream, or a WebSocket. + */ +type StreamClient = ServerResponse | WebSocketLikeClient; +/** + * One transport's clients. `createHot` publishes through this and does not know + * whether the events leave over Server-Sent Events, a WebSocket or something of + * your own, which is what `TClient` is for: a transport built by a `transport` + * function names the type of the clients it hands to `onConnect` and takes back + * in `publishTo`. + */ +type ClientStream = { + /** + * answer a request on the endpoint's path */ handler: (req: IncomingMessage, res: ServerResponse) => void; /** * true when at least one client is connected */ hasClients: () => boolean; + /** + * called with each client once it has joined + */ + onConnect: (fn: (client: TClient) => void) => void; /** * publish a payload to every client */ @@ -242,7 +275,7 @@ type EventStream = { * publish a payload to a single client */ publishTo: ( - res: ServerResponse, + client: TClient, payload: | Payload | { @@ -253,4 +286,24 @@ type EventStream = { * end every client and stop the heartbeat */ close: () => void; + /** + * answer upgrades on this server + */ + attach?: ((server: HttpServer) => void) | undefined; + /** + * stop answering upgrades + */ + detach?: (() => void) | undefined; }; +/** + * Builds a transport of your own. The same calls `createHot` makes of the + * built-in two are made of whatever this returns. + */ +type ClientStreamFactory = ( + options: { + path: string; + heartbeat: number; + }, + logger: Logger, +) => ClientStream; +type EventStream = ClientStream; diff --git a/types/index.d.ts b/types/index.d.ts index 197325df9..f05b81cc5 100644 --- a/types/index.d.ts +++ b/types/index.d.ts @@ -53,6 +53,7 @@ declare namespace wdm { GetFilenameFromUrl, WaitUntilValid, Invalidate, + Attach, Close, AdditionalMethods, API, @@ -341,6 +342,7 @@ type GetFilenameFromUrl = ( ) => Promise; type WaitUntilValid = (callback: Callback) => any; type Invalidate = (callback: Callback) => any; +type Attach = (server: import("node:http").Server) => any; type Close = (callback: (err: Error | null | undefined) => void) => any; type AdditionalMethods< RequestInternal extends IncomingMessage, @@ -358,6 +360,10 @@ type AdditionalMethods< * invalidate */ invalidate: Invalidate; + /** + * answer WebSocket upgrades on this server + */ + attach: Attach; /** * close */ diff --git a/types/servers/WebSocketServer.d.ts b/types/servers/WebSocketServer.d.ts new file mode 100644 index 000000000..aa1db3e9d --- /dev/null +++ b/types/servers/WebSocketServer.d.ts @@ -0,0 +1,52 @@ +export = createWebSocketStream; +/** + * A client stream carried over WebSocket rather than Server-Sent Events. It + * answers the same calls as `createEventStream`, so `createHot` does not know + * which of them it is publishing to. + * @param {object} options options + * @param {string} options.path the path the endpoint is served at + * @param {number} options.heartbeat heartbeat interval in milliseconds + * @param {Logger} logger logger + * @returns {ClientStream} client stream + */ +declare function createWebSocketStream( + { + path, + heartbeat, + }: { + path: string; + heartbeat: number; + }, + logger: Logger, +): ClientStream; +declare namespace createWebSocketStream { + export { + WS_DEFAULT_HEARTBEAT, + createWebSocketStream, + HttpServer, + IncomingMessage, + Duplex, + WebSocket, + WsServerConstructor, + Logger, + Payload, + ClientStream, + }; +} +/** @typedef {import("node:http").Server} HttpServer */ +/** @typedef {import("node:http").IncomingMessage} IncomingMessage */ +/** @typedef {import("node:stream").Duplex} Duplex */ +/** @typedef {import("ws").WebSocket} WebSocket */ +/** @typedef {typeof import("ws").WebSocketServer} WsServerConstructor */ +/** @typedef {import("../hot.js").Logger} Logger */ +/** @typedef {import("../hot.js").Payload} Payload */ +/** @typedef {import("../hot.js").ClientStream} ClientStream */ +declare const WS_DEFAULT_HEARTBEAT: number; +type HttpServer = import("node:http").Server; +type IncomingMessage = import("node:http").IncomingMessage; +type Duplex = import("node:stream").Duplex; +type WebSocket = import("ws").WebSocket; +type WsServerConstructor = typeof import("ws").WebSocketServer; +type Logger = import("../hot.js").Logger; +type Payload = import("../hot.js").Payload; +type ClientStream = import("../hot.js").ClientStream;