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

Function toLayerQueue

packages/effect/src/unstable/cluster/Entity.ts:310–400  ·  view source on GitHub ↗
(
    this: Entity<string, Rpcs>,
    build:
      | ((
        mailbox: Queue.Dequeue<Envelope.Request<Rpcs>>,
        replier: Replier<Rpcs>
      ) => Effect.Effect<never, never, R>)
      | Effect.Effect<
        (
          mailbox: Queue.Dequeue<Envelope.Request<Rpcs>>,
          replier: Replier<Rpcs>
        ) => Effect.Effect<never, never, R>,
        never,
        RX
      >,
    options?: {
      readonly maxIdleTime?: Duration.Input | undefined
      readonly mailboxCapacity?: number | "unbounded" | undefined
      readonly disableFatalDefects?: boolean | undefined
      readonly defectRetryPolicy?: Schedule.Schedule<any, unknown> | undefined
      readonly spanAttributes?: Record<string, string> | undefined
    }
  )

Source from the content-addressed store, hash-verified

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)

Callers

nothing calls this directly

Calls 7

completeFunction · 0.70
resumeFunction · 0.70
makeMethod · 0.65
pipeMethod · 0.65
toLayerMethod · 0.65
succeedMethod · 0.45
failMethod · 0.45

Tested by

no test coverage detected