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

Function fromWebSocket

packages/effect/src/unstable/socket/Socket.ts:606–760  ·  view source on GitHub ↗
(
  acquire: Effect.Effect<globalThis.WebSocket, SocketError, RO>,
  options?: {
    readonly closeCodeIsError?: ((code: number) => boolean) | undefined
    readonly openTimeout?: Duration.Input | undefined
    /**
     * Replays buffered events on the first run after the socket opens and before
     * the run's `onOpen` effect.
     *
     * @category options
     * @since 4.0.0
     */
    readonly onInitialRun?: ((ws: globalThis.WebSocket) => ReadonlyArray<MessageEvent>) | undefined
  } | undefined
)

Source from the content-addressed store, hash-verified

604 * @since 4.0.0
605 */
606export const fromWebSocket = <RO>(
607 acquire: Effect.Effect<globalThis.WebSocket, SocketError, RO>,
608 options?: {
609 readonly closeCodeIsError?: ((code: number) => boolean) | undefined
610 readonly openTimeout?: Duration.Input | undefined
611 /**
612 * Replays buffered events on the first run after the socket opens and before
613 * the run's `onOpen` effect.
614 *
615 * @category options
616 * @since 4.0.0
617 */
618 readonly onInitialRun?: ((ws: globalThis.WebSocket) => ReadonlyArray<MessageEvent>) | undefined
619 } | undefined
620): Effect.Effect<Socket, never, Exclude<RO, Scope.Scope>> =>
621 Effect.withFiber((fiber) => {
622 let currentWS: globalThis.WebSocket | undefined
623 let initial = true
624 const latch = Latch.makeUnsafe(false)
625 const acquireContext = fiber.context as Context.Context<RO>
626 const closeCodeIsError = options?.closeCodeIsError ?? defaultCloseCodeIsError
627
628 const runRaw = <_, E, R>(handler: (_: string | Uint8Array) => Effect.Effect<_, E, R> | void, opts?: {
629 readonly onOpen?: Effect.Effect<void> | undefined
630 }) =>
631 Effect.scopedWith(Effect.fnUntraced(function*(scope) {
632 const fiberSet = yield* FiberSet.make<any, E | SocketError>().pipe(
633 Scope.provide(scope)
634 )
635 const ws = yield* Scope.provide(acquire, scope)
636 const run = yield* Effect.provideService(FiberSet.runtime(fiberSet)<R>(), WebSocket, ws)
637 let open = false
638
639 function onMessage(event: MessageEvent) {
640 if (event.data instanceof Blob) {
641 const effect = Effect.flatMap(
642 Effect.promise(() => event.data.arrayBuffer() as Promise<ArrayBuffer>),
643 (buffer) => {
644 const result = handler(new Uint8Array(buffer))
645 return Effect.isEffect(result) ? result : Effect.void
646 }
647 )
648 return run(effect)
649 }
650 const result = handler(event.data instanceof ArrayBuffer ? new Uint8Array(event.data) : event.data)
651 if (Effect.isEffect(result)) {
652 run(result)
653 }
654 }
655 function onError(cause: Event) {
656 ws.removeEventListener("message", onMessage)
657 ws.removeEventListener("close", onClose)
658 Deferred.doneUnsafe(
659 fiberSet.deferred,
660 Effect.fail(
661 new SocketError({
662 reason: open ?
663 new SocketReadError({

Callers 1

makeWebSocketFunction · 0.85

Calls 2

makeFunction · 0.70
succeedMethod · 0.45

Tested by

no test coverage detected