(
requestFiber: Fiber.Fiber<any, any>,
client: Client,
request: Request<Rpcs>,
opts: Parameters<RpcServer<Rpcs>["write"]>[2]
)
| 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 { |
no test coverage detected