diff --git a/packages/engine.io/lib/server.ts b/packages/engine.io/lib/server.ts index c7c571613..2245900ba 100644 --- a/packages/engine.io/lib/server.ts +++ b/packages/engine.io/lib/server.ts @@ -685,11 +685,17 @@ export abstract class BaseServer extends EventEmitter { * * @see https://nodejs.org/api/http.html#class-httpserverresponse */ -class WebSocketResponse { +class WebSocketResponse extends EventEmitter { + /** + * The status code of the handshake response written by the "ws" package. + */ + public statusCode = 101; + constructor( readonly req, readonly socket: Duplex, ) { + super(); // temporarily store the response headers on the req object (see the "headers" event) req[kResponseHeaders] = {}; } @@ -875,6 +881,11 @@ export class Server extends BaseServer { // delegate to ws this.ws.handleUpgrade(engineRequest, socket, head, (websocket) => { + // the handshake response has been written, so middlewares which track the + // response lifecycle (loggers, metrics) can now complete + res.emit("finish"); + res.emit("close"); + this.onWebSocket(engineRequest, socket, websocket); }); }; diff --git a/packages/engine.io/test/middlewares.js b/packages/engine.io/test/middlewares.js index 60e1cc6b2..000f7a7e6 100644 --- a/packages/engine.io/test/middlewares.js +++ b/packages/engine.io/test/middlewares.js @@ -249,6 +249,38 @@ describe("middlewares", () => { }); }); + it("should support middlewares that track the response lifecycle (websocket)", function (done) { + // the uWebSockets.js server has its own response wrapper, see ResponseWrapper in lib/userver.ts + if (process.env.EIO_WS_ENGINE === "uws") { + return this.skip(); + } + const engine = listen((port) => { + let responseClosed = false; + + engine.use((req, res, next) => { + // response loggers such as pino-http register listeners on the response + res.on("close", () => { + responseClosed = true; + }); + next(); + }); + + const socket = new WebSocket( + `ws://localhost:${port}/engine.io/?EIO=4&transport=websocket`, + ); + + socket.on("open", () => { + expect(responseClosed).to.be(true); + + socket.close(); + if (engine.httpServer) { + engine.httpServer.close(); + } + done(); + }); + }); + }); + it("should fail on errors (polling)", (done) => { const engine = listen((port) => { engine.use((req, res, next) => {