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

Function handleRequest

packages/effect/src/unstable/rpc/RpcServer.ts:234–394  ·  view source on GitHub ↗
(
    requestFiber: Fiber.Fiber<any, any>,
    client: Client,
    request: Request<Rpcs>,
    opts: Parameters<RpcServer<Rpcs>["write"]>[2]
  )

Source from the content-addressed store, hash-verified

232 }
233
234 const handleRequest = (
235 requestFiber: Fiber.Fiber<any, any>,
236 client: Client,
237 request: Request<Rpcs>,
238 opts: Parameters<RpcServer<Rpcs>["write"]>[2]
239 ): Effect.Effect<void> => {
240 if (client.fibers.has(request.id)) {
241 return Effect.interrupt
242 }
243 const rpc = group.requests.get(request.tag) as any as Rpc.AnyWithProps
244 const entry = Context.getOrUndefinedUnsafe(services, rpc?.key) as Rpc.Handler<Rpcs["_tag"]>
245 if (!rpc || !entry) {
246 const write = Effect.catchDefect(
247 options.onFromServer({
248 _tag: "Exit",
249 clientId: client.id,
250 requestId: request.id,
251 exit: Exit.die(`Unknown request tag: ${request.tag}`)
252 }),
253 (defect) => sendDefect(client, defect)
254 )
255 if (!client.ended || client.fibers.size > 0) return write
256 return Effect.ensuring(write, endClient(client))
257 }
258 const isStream = RpcSchema.isStreamSchema(rpc.successSchema)
259 const metadata = {
260 rpc,
261 client: client.serverClient,
262 requestId: request.id,
263 headers: request.headers,
264 payload: request.payload
265 }
266 const result = entry.handler(request.payload, metadata)
267
268 // if the handler requested forking, then we skip the concurrency control
269 const isWrapper = Rpc.isWrapper(result)
270 const isFork = isWrapper && result.fork
271 const isUninterruptible = isWrapper && result.uninterruptible
272 // unwrap the fork data type
273 const streamOrEffect = isWrapper ? result.value : result
274 const handler = isStream
275 ? (streamEffect(client, request, streamOrEffect) as Effect.Effect<{} | Deferred.Deferred<any, any>>)
276 : (streamOrEffect as Effect.Effect<{} | Deferred.Deferred<any, any>>)
277
278 const withMiddleware = rpc.middlewares.size > 0
279 ? applyMiddleware(services, handler, metadata)
280 : handler
281 let responded = false
282 const scope = Scope.makeUnsafe()
283 let deferred: Deferred.Deferred<unknown, unknown> | undefined = undefined
284 let effect = Effect.onExit(withMiddleware, (exit) => {
285 responded = true
286 let write: Effect.Effect<void>
287 if (exit._tag === "Success") {
288 if (Deferred.isDeferred(exit.value)) {
289 deferred = exit.value
290 write = Effect.void
291 } else {

Callers 1

writeFunction · 0.70

Calls 14

sendDefectFunction · 0.85
endClientFunction · 0.85
streamEffectFunction · 0.85
runForkFunction · 0.85
withPermitMethod · 0.80
addObserverMethod · 0.80
applyMiddlewareFunction · 0.70
getMethod · 0.65
closeUnsafeMethod · 0.65
forEachMethod · 0.65
setMethod · 0.65
awaitMethod · 0.65

Tested by

no test coverage detected