(
lazyWss: Effect.Effect<NodeWS.WebSocketServer>,
httpEffect: Effect.Effect<HttpServerResponse, E, R>,
options: {
readonly scope: Scope.Scope
readonly middleware?: Middleware.HttpMiddleware.Applied<App, E, R> | undefined
}
)
| 226 | * @since 4.0.0 |
| 227 | */ |
| 228 | export const makeUpgradeHandler = < |
| 229 | R, |
| 230 | E, |
| 231 | App extends Effect.Effect<HttpServerResponse, any, any> = Effect.Effect<HttpServerResponse, E, R> |
| 232 | >( |
| 233 | lazyWss: Effect.Effect<NodeWS.WebSocketServer>, |
| 234 | httpEffect: Effect.Effect<HttpServerResponse, E, R>, |
| 235 | options: { |
| 236 | readonly scope: Scope.Scope |
| 237 | readonly middleware?: Middleware.HttpMiddleware.Applied<App, E, R> | undefined |
| 238 | } |
| 239 | ): Effect.Effect< |
| 240 | (nodeRequest: Http.IncomingMessage, socket: Duplex, head: Buffer) => void, |
| 241 | never, |
| 242 | Exclude<Effect.Services<App>, HttpServerRequest | Scope.Scope> |
| 243 | > => { |
| 244 | const handledApp = HttpEffect.toHandled(httpEffect, handleResponse, options.middleware as any) |
| 245 | return Effect.withFiber((parent) => { |
| 246 | const services = parent.context |
| 247 | return Effect.succeed(function handler( |
| 248 | nodeRequest: Http.IncomingMessage, |
| 249 | socket: Duplex, |
| 250 | head: Buffer |
| 251 | ) { |
| 252 | let nodeResponse_: Http.ServerResponse | undefined = undefined |
| 253 | const nodeResponse = () => { |
| 254 | if (nodeResponse_ === undefined) { |
| 255 | nodeResponse_ = new Http.ServerResponse(nodeRequest) |
| 256 | nodeResponse_.assignSocket(socket as any) |
| 257 | nodeResponse_.on("finish", () => { |
| 258 | socket.end() |
| 259 | }) |
| 260 | } |
| 261 | return nodeResponse_ |
| 262 | } |
| 263 | const upgradeEffect = Socket.fromWebSocket(Effect.flatMap( |
| 264 | lazyWss, |
| 265 | (wss) => |
| 266 | Effect.acquireRelease( |
| 267 | Effect.callback<globalThis.WebSocket>((resume) => |
| 268 | wss.handleUpgrade(nodeRequest, socket, head, (ws) => { |
| 269 | resume(Effect.succeed(ws as any)) |
| 270 | }) |
| 271 | ), |
| 272 | (ws) => Effect.sync(() => ws.close()) |
| 273 | ) |
| 274 | )) |
| 275 | const context = Context.add( |
| 276 | services, |
| 277 | HttpServerRequest, |
| 278 | new ServerRequestImpl(nodeRequest, nodeResponse, upgradeEffect) |
| 279 | ) |
| 280 | const fiber = Fiber.runIn(Effect.runForkWith(context as Context.Context<any>)(handledApp), options.scope) |
| 281 | socket.on("error", () => {}) |
| 282 | socket.on("close", () => { |
| 283 | if (!socket.writableEnded) { |
| 284 | fiber.interruptUnsafe(parent.id, ClientAbort.annotation) |
| 285 | } |
no test coverage detected