diff --git a/.changeset/rpc-socket-ping-timeout-hook.md b/.changeset/rpc-socket-ping-timeout-hook.md new file mode 100644 index 00000000000..18f39363b81 --- /dev/null +++ b/.changeset/rpc-socket-ping-timeout-hook.md @@ -0,0 +1,5 @@ +--- +"effect": patch +--- + +Add `onPingTimeout` to `RpcClient.makeProtocolSocket` and `RpcClient.layerProtocolSocket` to distinguish ping timeouts from other socket failures. It runs only when a ping timeout drops an open connection, before `onDisconnect` and in-flight call failures; defects are logged and ignored. diff --git a/packages/effect/src/rpc/RpcClient.ts b/packages/effect/src/rpc/RpcClient.ts index 615d168adc5..d6d9bb330c2 100644 --- a/packages/effect/src/rpc/RpcClient.ts +++ b/packages/effect/src/rpc/RpcClient.ts @@ -1050,11 +1050,19 @@ export const makeProtocolSocket = (options?: { readonly retryPolicy?: Schedule.Schedule | undefined /** * Runs for each retried `SocketOpenError` when `retryTransientErrors` is enabled. - * A ping timeout fails in-flight calls and is not reported through this hook. - * The returned `Effect` cannot fail with a typed error or require - * services; defects are logged and ignored so retries can continue. + * A ping timeout fails in-flight calls and is reported through `onPingTimeout` + * instead. The returned `Effect` cannot fail with a typed error or + * require services; defects are logged and ignored so retries can continue. */ readonly onTransientError?: ((error: RpcClientError) => Effect.Effect) | undefined + /** + * Runs when an open connection is dropped because no server frame arrived + * within `pingTimeout`, before `ConnectionHooks.onDisconnect` and before + * in-flight calls fail with a `SocketReadError`. Defects are logged and ignored. + * + * @since 4.0.3 + */ + readonly onPingTimeout?: Effect.Effect | undefined }): Effect.Effect< Protocol["Service"], never, @@ -1070,10 +1078,12 @@ export const makeProtocolSocket = (options?: { let parser = serialization.makeUnsafe() - // `parser` is replaced on every connect, and a stateful serialization - // encodes against the connection it is writing to, so the ping is encoded - // when it is sent rather than once up front. + // Encode each ping with the current connection's parser. const pinger = yield* makePinger(Effect.suspend(() => writer.write(parser.encode(constPing)!)), options) + const onPingTimeout = Effect.ignoreCause(options?.onPingTimeout ?? Effect.void, { + log: true, + message: "RpcClient onPingTimeout hook failed" + }) let currentError: RpcClientError | undefined const broadcast = (response: FromServerEncoded) => @@ -1127,9 +1137,11 @@ export const makeProtocolSocket = (options?: { yield* Effect.suspend(() => { parser = serialization.makeUnsafe() pinger.reset() + let connected = false return Effect.gen(function*() { const { pull } = yield* socket.reader currentError = undefined + connected = true if (Option.isSome(hooks)) { yield* hooks.value.onConnect } @@ -1141,17 +1153,19 @@ export const makeProtocolSocket = (options?: { } }).pipe( Effect.scoped, - Effect.raceFirst(Effect.flatMap( - pinger.timeout, - () => - Effect.fail( - new Socket.SocketError({ - reason: new Socket.SocketReadError({ - cause: new Error("ping timeout") - }) + // The read loop never succeeds, so the race only succeeds on a ping + // timeout, after the socket has been cleaned up. + Effect.raceFirst(pinger.timeout), + Effect.andThen(() => connected ? onPingTimeout : Effect.void), + Effect.andThen(() => + Effect.fail( + new Socket.SocketError({ + reason: new Socket.SocketReadError({ + cause: new Error("ping timeout") }) - ) - )) + }) + ) + ) ) }).pipe( Option.isSome(hooks) ? Effect.ensuring(hooks.value.onDisconnect) : identity, @@ -1262,11 +1276,19 @@ export const layerProtocolSocket = (options?: { readonly retryPolicy?: Schedule.Schedule | undefined /** * Runs for each retried `SocketOpenError` when `retryTransientErrors` is enabled. - * A ping timeout fails in-flight calls and is not reported through this hook. - * The returned `Effect` cannot fail with a typed error or require - * services; defects are logged and ignored so retries can continue. + * A ping timeout fails in-flight calls and is reported through `onPingTimeout` + * instead. The returned `Effect` cannot fail with a typed error or + * require services; defects are logged and ignored so retries can continue. */ readonly onTransientError?: ((error: RpcClientError) => Effect.Effect) | undefined + /** + * Runs when an open connection is dropped because no server frame arrived + * within `pingTimeout`, before `ConnectionHooks.onDisconnect` and before + * in-flight calls fail with a `SocketReadError`. Defects are logged and ignored. + * + * @since 4.0.3 + */ + readonly onPingTimeout?: Effect.Effect | undefined }): Layer.Layer< Protocol, never, diff --git a/packages/effect/test/rpc/RpcClient.test.ts b/packages/effect/test/rpc/RpcClient.test.ts index 88c13947dd0..78c2c4962e6 100644 --- a/packages/effect/test/rpc/RpcClient.test.ts +++ b/packages/effect/test/rpc/RpcClient.test.ts @@ -510,6 +510,101 @@ describe("RpcClient", () => { assert.strictEqual(error.reason._tag, "SocketReadError") })) + it.effect("runs onPingTimeout only when a ping timeout drops the socket", () => + Effect.gen(function*() { + const requestSent = yield* Deferred.make() + const events: Array = [] + const record = (event: string) => Effect.sync(() => events.push(event)) + const closeError = new Socket.SocketError({ + reason: new Socket.SocketCloseError({ code: 1006 }) + }) + let connections = 0 + const write = () => Effect.asVoid(Deferred.succeed(requestSent, void 0)) + const socket = Socket.make({ + // The server closes the first connection; the reconnected one goes silent. + reader: Effect.sync(() => ({ + pull: ++connections === 1 + ? Deferred.await(requestSent).pipe(Effect.andThen(Effect.fail(closeError))) + : Effect.never, + upgrade: Socket.SocketUpgradeError.unsupported + })), + writer: Effect.succeed({ write, writeAll: write }) + }) + const protocol = yield* RpcClient.makeProtocolSocket({ + retryPolicy: Schedule.spaced("20 seconds"), + onPingTimeout: record("ping timeout") + }).pipe( + Effect.provideService(Socket.Socket, socket), + Effect.provideService(RpcClient.ConnectionHooks, { + onConnect: Effect.asVoid(record("connect")), + onDisconnect: Effect.asVoid(record("disconnect")) + }), + Effect.provide(RpcSerialization.layerNdjson) + ) + const client = yield* RpcClient.make(TestGroup).pipe( + Effect.provideService(RpcClient.Protocol, protocol) + ) + + const closedFiber = yield* client.Events().pipe(Stream.runDrain, Effect.flip, Effect.forkChild) + const closed = yield* Fiber.join(closedFiber) + assert.strictEqual(closed.reason._tag, "SocketCloseError") + + // Stay disconnected past the ping timeout, then reconnect at 20s. + yield* TestClock.adjust("20 seconds") + assert.deepStrictEqual(events, ["connect", "disconnect", "connect"]) + + const streamFiber = yield* client.Events().pipe( + Stream.runDrain, + Effect.tapError(() => record("stream failed")), + Effect.flip, + Effect.forkChild + ) + yield* TestClock.adjust("10 seconds") + const timedOut = yield* Fiber.join(streamFiber) + + assert.strictEqual(timedOut.reason._tag, "SocketReadError") + assert.deepStrictEqual(events, [ + "connect", + "disconnect", + "connect", + "ping timeout", + "disconnect", + "stream failed" + ]) + })) + + it.effect("does not run onPingTimeout when the ping timeout fires before the socket opens", () => + Effect.gen(function*() { + const events: Array = [] + const record = (event: string) => Effect.sync(() => events.push(event)) + const write = () => Effect.void + const socket = Socket.make({ + reader: Effect.never, + writer: Effect.succeed({ write, writeAll: write }) + }) + const protocol = yield* RpcClient.makeProtocolSocket({ + retryPolicy: Schedule.spaced("1 hour"), + onPingTimeout: record("ping timeout") + }).pipe( + Effect.provideService(Socket.Socket, socket), + Effect.provideService(RpcClient.ConnectionHooks, { + onConnect: Effect.asVoid(record("connect")), + onDisconnect: Effect.asVoid(record("disconnect")) + }), + Effect.provide(RpcSerialization.layerNdjson) + ) + const client = yield* RpcClient.make(TestGroup).pipe( + Effect.provideService(RpcClient.Protocol, protocol) + ) + const streamFiber = yield* client.Events().pipe(Stream.runDrain, Effect.flip, Effect.forkChild) + + yield* TestClock.adjust("11 seconds") + const error = yield* Fiber.join(streamFiber) + + assert.strictEqual(error.reason._tag, "SocketReadError") + assert.deepStrictEqual(events, ["disconnect"]) + })) + it.effect("fails in-flight streams when transient retries are exhausted", () => Effect.gen(function*() { const requestSent = yield* Deferred.make()