MCPcopy Create free account
hub / github.com/Effect-TS/effect / makeUpgradeHandler

Function makeUpgradeHandler

packages/platform/node/src/NodeHttpServer.ts:228–289  ·  view source on GitHub ↗
(
  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
  }
)

Source from the content-addressed store, hash-verified

226 * @since 4.0.0
227 */
228export 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 }

Callers 1

NodeHttpServer.tsFile · 0.85

Calls 7

interruptUnsafeMethod · 0.80
closeMethod · 0.65
addMethod · 0.65
onMethod · 0.65
resumeFunction · 0.50
succeedMethod · 0.45
syncMethod · 0.45

Tested by

no test coverage detected