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

Function runRaw

packages/effect/src/unstable/socket/Socket.ts:628–736  ·  view source on GitHub ↗
(handler: (_: string | Uint8Array) => Effect.Effect<_, E, R> | void, opts?: {
      readonly onOpen?: Effect.Effect<void> | undefined
    })

Source from the content-addressed store, hash-verified

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({
664 cause
665 }) :
666 new SocketOpenError({
667 kind: "Unknown",
668 cause
669 })
670 })
671 )
672 )
673 }
674 function onClose(event: globalThis.CloseEvent) {
675 const code = typeof event.code === "number" ? event.code : 1001
676 ws.removeEventListener("message", onMessage)
677 ws.removeEventListener("error", onError)
678 Deferred.doneUnsafe(
679 fiberSet.deferred,
680 Effect.fail(
681 new SocketError({
682 reason: new SocketCloseError({
683 code,
684 closeReason: event.reason
685 })

Callers

nothing calls this directly

Calls 15

runForkFunction · 0.85
isSocketErrorFunction · 0.85
joinMethod · 0.80
filterCleanMethod · 0.80
mergeMethod · 0.80
addFinalizerMethod · 0.80
onMessageFunction · 0.70
pipeMethod · 0.65
makeMethod · 0.65
addEventListenerMethod · 0.65
awaitMethod · 0.65
openUnsafeMethod · 0.65

Tested by

no test coverage detected