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

Function makeProtocolSocket

packages/effect/src/unstable/rpc/RpcClient.ts:1011–1161  ·  view source on GitHub ↗
(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
})

Source from the content-addressed store, hash-verified

1009 * @since 4.0.0
1010 */
1011export 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 }

Callers 1

layerProtocolSocketFunction · 0.85

Calls 9

broadcastFunction · 0.85
broadcastErrorFunction · 0.85
writeFunction · 0.70
makeMethod · 0.65
pipeMethod · 0.65
resetMethod · 0.65
getMethod · 0.65
runRawMethod · 0.45
failMethod · 0.45

Tested by

no test coverage detected