(options?: {
readonly retryTransientErrors?: boolean | undefined
readonly retryPolicy?: Schedule.Schedule<any, Socket.SocketError> | undefined
/**
* Runs for each retried `SocketOpenError` when `retryTransientErrors` is enabled.
* Ping timeouts are also reported because the protocol classifies them as
* `SocketOpenError`. The returned `Effect<void>` cannot fail with a typed error
* or require services; defects are logged and ignored so retries can continue.
*/
readonly onTransientError?: ((error: RpcClientError) => Effect.Effect<void>) | undefined
})
| 1009 | * @since 4.0.0 |
| 1010 | */ |
| 1011 | export const makeProtocolSocket = (options?: { |
| 1012 | readonly retryTransientErrors?: boolean | undefined |
| 1013 | readonly retryPolicy?: Schedule.Schedule<any, Socket.SocketError> | undefined |
| 1014 | /** |
| 1015 | * Runs for each retried `SocketOpenError` when `retryTransientErrors` is enabled. |
| 1016 | * Ping timeouts are also reported because the protocol classifies them as |
| 1017 | * `SocketOpenError`. The returned `Effect<void>` cannot fail with a typed error |
| 1018 | * or require services; defects are logged and ignored so retries can continue. |
| 1019 | */ |
| 1020 | readonly onTransientError?: ((error: RpcClientError) => Effect.Effect<void>) | undefined |
| 1021 | }): Effect.Effect< |
| 1022 | Protocol["Service"], |
| 1023 | never, |
| 1024 | Scope.Scope | RpcSerialization.RpcSerialization | Socket.Socket |
| 1025 | > => |
| 1026 | Protocol.make(Effect.fnUntraced(function*(writeResponse, clientIds) { |
| 1027 | const socket = yield* Socket.Socket |
| 1028 | const serialization = yield* RpcSerialization.RpcSerialization |
| 1029 | const hooks = yield* Effect.serviceOption(ConnectionHooks) |
| 1030 | const requestClientMap = new Map<string | number, number>() |
| 1031 | |
| 1032 | const write = yield* socket.writer |
| 1033 | |
| 1034 | let parser = serialization.makeUnsafe() |
| 1035 | |
| 1036 | const pinger = yield* makePinger(write(parser.encode(constPing)!)) |
| 1037 | let currentError: RpcClientError | undefined |
| 1038 | const onOpen = Effect.suspend(() => { |
| 1039 | currentError = undefined |
| 1040 | return Option.isSome(hooks) ? hooks.value.onConnect : Effect.void |
| 1041 | }) |
| 1042 | |
| 1043 | const broadcast = (response: FromServerEncoded) => |
| 1044 | Effect.forEach(clientIds, (clientId) => writeResponse(clientId, response)) |
| 1045 | const broadcastError = (error: RpcClientError) => { |
| 1046 | currentError = error |
| 1047 | return broadcast({ |
| 1048 | _tag: "ClientProtocolError", |
| 1049 | error |
| 1050 | }) |
| 1051 | } |
| 1052 | |
| 1053 | yield* Effect.suspend(() => { |
| 1054 | parser = serialization.makeUnsafe() |
| 1055 | pinger.reset() |
| 1056 | return socket.runRaw((message) => { |
| 1057 | try { |
| 1058 | const responses = parser.decode(message) as Array<FromServerEncoded> |
| 1059 | if (responses.length === 0) return |
| 1060 | let i = 0 |
| 1061 | return Effect.whileLoop({ |
| 1062 | while: () => i < responses.length, |
| 1063 | body: () => { |
| 1064 | const response = responses[i++] |
| 1065 | if (response._tag === "Pong") { |
| 1066 | pinger.onPong() |
| 1067 | return Effect.void |
| 1068 | } |
no test coverage detected