(
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
)
| 604 | * @since 4.0.0 |
| 605 | */ |
| 606 | export 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({ |
no test coverage detected