| 308 | }, |
| 309 | of: identity, |
| 310 | toLayerQueue< |
| 311 | Rpcs extends Rpc.Any, |
| 312 | R, |
| 313 | RX = never |
| 314 | >( |
| 315 | this: Entity<string, Rpcs>, |
| 316 | build: |
| 317 | | (( |
| 318 | mailbox: Queue.Dequeue<Envelope.Request<Rpcs>>, |
| 319 | replier: Replier<Rpcs> |
| 320 | ) => Effect.Effect<never, never, R>) |
| 321 | | Effect.Effect< |
| 322 | ( |
| 323 | mailbox: Queue.Dequeue<Envelope.Request<Rpcs>>, |
| 324 | replier: Replier<Rpcs> |
| 325 | ) => Effect.Effect<never, never, R>, |
| 326 | never, |
| 327 | RX |
| 328 | >, |
| 329 | options?: { |
| 330 | readonly maxIdleTime?: Duration.Input | undefined |
| 331 | readonly mailboxCapacity?: number | "unbounded" | undefined |
| 332 | readonly disableFatalDefects?: boolean | undefined |
| 333 | readonly defectRetryPolicy?: Schedule.Schedule<any, unknown> | undefined |
| 334 | readonly spanAttributes?: Record<string, string> | undefined |
| 335 | } |
| 336 | ) { |
| 337 | const buildHandlers = Effect.gen({ self: this }, function*() { |
| 338 | const behaviour = Effect.isEffect(build) ? yield* build : build |
| 339 | const queue = yield* Queue.make<Envelope.Request<Rpcs>>() |
| 340 | |
| 341 | // create the rpc handlers for the entity |
| 342 | const handler = (envelope: any) => |
| 343 | Effect.callback<any, any>((resume) => { |
| 344 | Queue.offerUnsafe(queue, envelope) |
| 345 | resumes.set(envelope, resume) |
| 346 | }) |
| 347 | const streamHandler = (envelope: any) => |
| 348 | Effect.callback<any, any>((resume) => { |
| 349 | Queue.offerUnsafe(queue, envelope) |
| 350 | resumes.set(envelope, resume) |
| 351 | }).pipe( |
| 352 | Effect.map((streamOrQueue) => |
| 353 | Stream.isStream(streamOrQueue) ? streamOrQueue : Stream.fromQueue(streamOrQueue) |
| 354 | ), |
| 355 | Stream.unwrap |
| 356 | ) |
| 357 | const handlers: Record<string, any> = Object.create(null) |
| 358 | for (const rpc_ of this.protocol.requests.values()) { |
| 359 | const rpc = rpc_ as any as Rpc.AnyWithProps |
| 360 | handlers[rpc._tag] = RpcSchema.isStreamSchema(rpc.successSchema) ? streamHandler : handler |
| 361 | } |
| 362 | |
| 363 | // make the Replier for the behaviour |
| 364 | const resumes = new Map<Envelope.Request<any>, (exit: Exit.Exit<any, any>) => void>() |
| 365 | const complete = (request: Envelope.Request<any>, exit: Exit.Exit<any, any>) => |
| 366 | Effect.sync(() => { |
| 367 | const resume = resumes.get(request) |