From e94a0c9b4da3196b6e61b254686fc1eccddee5a4 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Fri, 9 Oct 2026 12:42:02 -0700 Subject: [PATCH 1/6] feat(rpc): report socket ping timeouts through an onPingTimeout option Add `onPingTimeout` to `RpcClient.makeProtocolSocket` and `layerProtocolSocket`. It runs when the pinger drops the connection because no server frame arrived within `pingTimeout`, before `ConnectionHooks.onDisconnect` and before in-flight calls fail with `SocketReadError`. Defects are logged and ignored, like `onTransientError`. Co-Authored-By: Claude Opus 5.5 --- .changeset/rpc-socket-ping-timeout-hook.md | 5 + packages/effect/src/rpc/RpcClient.ts | 53 ++++++-- packages/effect/test/rpc/RpcClient.test.ts | 146 +++++++++++++++++++++ 3 files changed, 192 insertions(+), 12 deletions(-) create mode 100644 .changeset/rpc-socket-ping-timeout-hook.md diff --git a/.changeset/rpc-socket-ping-timeout-hook.md b/.changeset/rpc-socket-ping-timeout-hook.md new file mode 100644 index 00000000000..78e70559060 --- /dev/null +++ b/.changeset/rpc-socket-ping-timeout-hook.md @@ -0,0 +1,5 @@ +--- +"effect": patch +--- + +Add an `onPingTimeout` option to `RpcClient.makeProtocolSocket` and `RpcClient.layerProtocolSocket`. It runs when a ping timeout drops the connection, before `ConnectionHooks.onDisconnect`, so clients can tell a server that stopped responding from one that closed the connection. diff --git a/packages/effect/src/rpc/RpcClient.ts b/packages/effect/src/rpc/RpcClient.ts index 615d168adc5..a9b674d4c66 100644 --- a/packages/effect/src/rpc/RpcClient.ts +++ b/packages/effect/src/rpc/RpcClient.ts @@ -1038,6 +1038,8 @@ export const layerProtocolHttp = (options: { * `pingInterval` defaults to 5 seconds. `pingTimeout` defaults to the interval * and measures time since the last decoded server frame, not the last pong. * The timeout is checked on each ping tick; any decoded frame counts as liveness. + * When the timeout drops the connection, `onPingTimeout` runs before + * `ConnectionHooks.onDisconnect`. * * @stability unstable * @category protocols @@ -1050,11 +1052,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 the 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, @@ -1074,6 +1084,12 @@ export const makeProtocolSocket = (options?: { // encodes against the connection it is writing to, so the ping is encoded // when it is sent rather than once up front. const pinger = yield* makePinger(Effect.suspend(() => writer.write(parser.encode(constPing)!)), options) + const onPingTimeout = options?.onPingTimeout ? + Effect.ignoreCause(options.onPingTimeout, { + log: true, + message: "RpcClient onPingTimeout hook failed" + }) : + Effect.void let currentError: RpcClientError | undefined const broadcast = (response: FromServerEncoded) => @@ -1144,12 +1160,15 @@ export const makeProtocolSocket = (options?: { Effect.raceFirst(Effect.flatMap( pinger.timeout, () => - Effect.fail( - new Socket.SocketError({ - reason: new Socket.SocketReadError({ - cause: new Error("ping timeout") + Effect.andThen( + onPingTimeout, + Effect.fail( + new Socket.SocketError({ + reason: new Socket.SocketReadError({ + cause: new Error("ping timeout") + }) }) - }) + ) ) )) ) @@ -1249,7 +1268,9 @@ const makePinger = Effect.fnUntraced(function*(writePing: Effect.Effect * `pingInterval` defaults to 5 seconds. `pingTimeout` defaults to the interval * and measures time since the last decoded server frame, not the last pong. * The timeout is checked on each ping tick; any decoded frame counts as liveness. - * `retryPolicy` configures retries after socket errors. + * When the timeout drops the connection, `onPingTimeout` runs before + * `ConnectionHooks.onDisconnect`. `retryPolicy` configures retries after socket + * errors. * * @stability unstable * @category layers @@ -1262,11 +1283,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 the 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..f5ad1e308c8 100644 --- a/packages/effect/test/rpc/RpcClient.test.ts +++ b/packages/effect/test/rpc/RpcClient.test.ts @@ -510,6 +510,152 @@ describe("RpcClient", () => { assert.strictEqual(error.reason._tag, "SocketReadError") })) + it.effect("runs onPingTimeout before onDisconnect 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 write = () => Effect.asVoid(Deferred.succeed(requestSent, void 0)) + const socket = Socket.make({ + reader: Effect.succeed({ + pull: Effect.never, + upgrade: Socket.SocketUpgradeError.unsupported + }), + 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.tapError(() => record("stream failed")), + Effect.flip, + Effect.forkChild + ) + + yield* Deferred.await(requestSent) + yield* TestClock.adjust("11 seconds") + const error = yield* Fiber.join(streamFiber) + + assert.instanceOf(error, RpcClientError) + assert.strictEqual(error.reason._tag, "SocketReadError") + assert.deepStrictEqual(events, ["connect", "ping timeout", "disconnect", "stream failed"]) + })) + + it.effect("does not run onPingTimeout when the socket closes for another reason", () => + 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 streamFiber = yield* client.Events().pipe(Stream.runDrain, Effect.flip, Effect.forkChild) + + const error = yield* Fiber.join(streamFiber) + assert.instanceOf(error, RpcClientError) + assert.strictEqual(error.reason._tag, "SocketCloseError") + assert.deepStrictEqual(events, ["connect", "disconnect"]) + + // Stay disconnected well past the ping timeout. + yield* TestClock.adjust("19 seconds") + assert.deepStrictEqual(events, ["connect", "disconnect"]) + + // Reconnect at 20s; the silent socket times out at the 30s tick. + yield* TestClock.adjust("12 seconds") + assert.deepStrictEqual(events, ["connect", "disconnect", "connect", "ping timeout", "disconnect"]) + })) + + it.effect("does not run onPingTimeout while server frames keep arriving", () => + Effect.gen(function*() { + const requestSent = yield* Deferred.make() + const frames = yield* Queue.unbounded() + const chunks = yield* Queue.unbounded() + const events: Array = [] + const record = (event: string) => Effect.sync(() => events.push(event)) + const write = () => Effect.asVoid(Deferred.succeed(requestSent, void 0)) + const socket = Socket.make({ + reader: Effect.succeed({ + pull: Queue.take(frames).pipe(Effect.map((frame) => [frame] as const)), + upgrade: Socket.SocketUpgradeError.unsupported + }), + 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, { + generateRequestId: () => RpcMessage.RequestId("0") + }).pipe(Effect.provideService(RpcClient.Protocol, protocol)) + const streamFiber = yield* client.Events().pipe( + Stream.runForEach((chunk) => Queue.offer(chunks, chunk)), + Effect.exit, + Effect.forkChild + ) + + yield* Deferred.await(requestSent) + for (const chunk of ["first", "second", "third"]) { + yield* TestClock.adjust("4 seconds") + yield* Queue.offer(frames, JSON.stringify({ _tag: "Chunk", requestId: "0", values: [chunk] }) + "\n") + assert.strictEqual(yield* Queue.take(chunks), chunk) + } + // Every ping tick so far (5s, 10s, 15s) came within 5s of a frame. + yield* TestClock.adjust("4 seconds") + assert.deepStrictEqual(events, ["connect"]) + assert.isUndefined(streamFiber.pollUnsafe()) + + // Once frames stop, the 20s tick is 8s after the last frame. + yield* TestClock.adjust("5 seconds") + assert.deepStrictEqual(events, ["connect", "ping timeout", "disconnect"]) + assert.isDefined(streamFiber.pollUnsafe()) + })) + it.effect("fails in-flight streams when transient retries are exhausted", () => Effect.gen(function*() { const requestSent = yield* Deferred.make() From 7fd7594506800d746c2329c34ada550e8004e492 Mon Sep 17 00:00:00 2001 From: Effect Bot Date: Fri, 9 Oct 2026 20:10:50 +0000 Subject: [PATCH 2/6] test(rpc): reduce onPingTimeout coverage to one focused test --- packages/effect/test/rpc/RpcClient.test.ts | 127 ++++----------------- 1 file changed, 22 insertions(+), 105 deletions(-) diff --git a/packages/effect/test/rpc/RpcClient.test.ts b/packages/effect/test/rpc/RpcClient.test.ts index f5ad1e308c8..cc219502fdf 100644 --- a/packages/effect/test/rpc/RpcClient.test.ts +++ b/packages/effect/test/rpc/RpcClient.test.ts @@ -510,50 +510,7 @@ describe("RpcClient", () => { assert.strictEqual(error.reason._tag, "SocketReadError") })) - it.effect("runs onPingTimeout before onDisconnect 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 write = () => Effect.asVoid(Deferred.succeed(requestSent, void 0)) - const socket = Socket.make({ - reader: Effect.succeed({ - pull: Effect.never, - upgrade: Socket.SocketUpgradeError.unsupported - }), - 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.tapError(() => record("stream failed")), - Effect.flip, - Effect.forkChild - ) - - yield* Deferred.await(requestSent) - yield* TestClock.adjust("11 seconds") - const error = yield* Fiber.join(streamFiber) - - assert.instanceOf(error, RpcClientError) - assert.strictEqual(error.reason._tag, "SocketReadError") - assert.deepStrictEqual(events, ["connect", "ping timeout", "disconnect", "stream failed"]) - })) - - it.effect("does not run onPingTimeout when the socket closes for another reason", () => + it.effect("runs onPingTimeout only when a ping timeout drops the socket", () => Effect.gen(function*() { const requestSent = yield* Deferred.make() const events: Array = [] @@ -588,72 +545,32 @@ describe("RpcClient", () => { Effect.provideService(RpcClient.Protocol, protocol) ) - const streamFiber = yield* client.Events().pipe(Stream.runDrain, Effect.flip, Effect.forkChild) - - const error = yield* Fiber.join(streamFiber) - assert.instanceOf(error, RpcClientError) - assert.strictEqual(error.reason._tag, "SocketCloseError") - assert.deepStrictEqual(events, ["connect", "disconnect"]) - - // Stay disconnected well past the ping timeout. - yield* TestClock.adjust("19 seconds") - assert.deepStrictEqual(events, ["connect", "disconnect"]) + const closedFiber = yield* client.Events().pipe(Stream.runDrain, Effect.flip, Effect.forkChild) + const closed = yield* Fiber.join(closedFiber) + assert.strictEqual(closed.reason._tag, "SocketCloseError") - // Reconnect at 20s; the silent socket times out at the 30s tick. - yield* TestClock.adjust("12 seconds") - assert.deepStrictEqual(events, ["connect", "disconnect", "connect", "ping timeout", "disconnect"]) - })) + // Stay disconnected past the ping timeout, then reconnect at 20s. + yield* TestClock.adjust("20 seconds") + assert.deepStrictEqual(events, ["connect", "disconnect", "connect"]) - it.effect("does not run onPingTimeout while server frames keep arriving", () => - Effect.gen(function*() { - const requestSent = yield* Deferred.make() - const frames = yield* Queue.unbounded() - const chunks = yield* Queue.unbounded() - const events: Array = [] - const record = (event: string) => Effect.sync(() => events.push(event)) - const write = () => Effect.asVoid(Deferred.succeed(requestSent, void 0)) - const socket = Socket.make({ - reader: Effect.succeed({ - pull: Queue.take(frames).pipe(Effect.map((frame) => [frame] as const)), - upgrade: Socket.SocketUpgradeError.unsupported - }), - 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, { - generateRequestId: () => RpcMessage.RequestId("0") - }).pipe(Effect.provideService(RpcClient.Protocol, protocol)) const streamFiber = yield* client.Events().pipe( - Stream.runForEach((chunk) => Queue.offer(chunks, chunk)), - Effect.exit, + Stream.runDrain, + Effect.tapError(() => record("stream failed")), + Effect.flip, Effect.forkChild ) - - yield* Deferred.await(requestSent) - for (const chunk of ["first", "second", "third"]) { - yield* TestClock.adjust("4 seconds") - yield* Queue.offer(frames, JSON.stringify({ _tag: "Chunk", requestId: "0", values: [chunk] }) + "\n") - assert.strictEqual(yield* Queue.take(chunks), chunk) - } - // Every ping tick so far (5s, 10s, 15s) came within 5s of a frame. - yield* TestClock.adjust("4 seconds") - assert.deepStrictEqual(events, ["connect"]) - assert.isUndefined(streamFiber.pollUnsafe()) - - // Once frames stop, the 20s tick is 8s after the last frame. - yield* TestClock.adjust("5 seconds") - assert.deepStrictEqual(events, ["connect", "ping timeout", "disconnect"]) - assert.isDefined(streamFiber.pollUnsafe()) + 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("fails in-flight streams when transient retries are exhausted", () => From 97318753f99a5f2e8513f8bd24ba1811805fdecf Mon Sep 17 00:00:00 2001 From: Effect Bot Date: Fri, 9 Oct 2026 20:20:48 +0000 Subject: [PATCH 3/6] fix(rpc): run onPingTimeout after the socket closes and only for open connections --- packages/effect/src/rpc/RpcClient.ts | 31 +++++++++++++++------------- 1 file changed, 17 insertions(+), 14 deletions(-) diff --git a/packages/effect/src/rpc/RpcClient.ts b/packages/effect/src/rpc/RpcClient.ts index a9b674d4c66..5e672bb8cca 100644 --- a/packages/effect/src/rpc/RpcClient.ts +++ b/packages/effect/src/rpc/RpcClient.ts @@ -1058,9 +1058,9 @@ export const makeProtocolSocket = (options?: { */ readonly onTransientError?: ((error: RpcClientError) => Effect.Effect) | undefined /** - * Runs when the 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. + * 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 */ @@ -1091,6 +1091,8 @@ export const makeProtocolSocket = (options?: { }) : Effect.void let currentError: RpcClientError | undefined + let connected = false + let pingTimeoutError: Socket.SocketError | undefined const broadcast = (response: FromServerEncoded) => Effect.forEach(clientIds, (clientId) => writeResponse(clientId, response)) @@ -1143,9 +1145,11 @@ export const makeProtocolSocket = (options?: { yield* Effect.suspend(() => { parser = serialization.makeUnsafe() pinger.reset() + connected = false return Effect.gen(function*() { const { pull } = yield* socket.reader currentError = undefined + connected = true if (Option.isSome(hooks)) { yield* hooks.value.onConnect } @@ -1160,19 +1164,18 @@ export const makeProtocolSocket = (options?: { Effect.raceFirst(Effect.flatMap( pinger.timeout, () => - Effect.andThen( - onPingTimeout, - Effect.fail( - new Socket.SocketError({ - reason: new Socket.SocketReadError({ - cause: new Error("ping timeout") - }) + Effect.fail( + pingTimeoutError = new Socket.SocketError({ + reason: new Socket.SocketReadError({ + cause: new Error("ping timeout") }) - ) + }) ) )) ) }).pipe( + // runs once the socket is closed, so a concurrent socket error cannot cut it short + Effect.tapError((error) => connected && error === pingTimeoutError ? onPingTimeout : Effect.void), Option.isSome(hooks) ? Effect.ensuring(hooks.value.onDisconnect) : identity, Effect.tapCause((cause) => { const error = Cause.findError(cause) @@ -1289,9 +1292,9 @@ export const layerProtocolSocket = (options?: { */ readonly onTransientError?: ((error: RpcClientError) => Effect.Effect) | undefined /** - * Runs when the 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. + * 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 */ From 7b83cf9db7004c0b646225b20725eee3c41b10d2 Mon Sep 17 00:00:00 2001 From: Effect Bot Date: Fri, 9 Oct 2026 20:23:38 +0000 Subject: [PATCH 4/6] test(rpc): cover ping timeouts that fire before the socket opens --- packages/effect/test/rpc/RpcClient.test.ts | 32 ++++++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/packages/effect/test/rpc/RpcClient.test.ts b/packages/effect/test/rpc/RpcClient.test.ts index cc219502fdf..78c2c4962e6 100644 --- a/packages/effect/test/rpc/RpcClient.test.ts +++ b/packages/effect/test/rpc/RpcClient.test.ts @@ -573,6 +573,38 @@ describe("RpcClient", () => { ]) })) + 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() From c7c53bf5967602b78c5dc8221adef1e5a68d5121 Mon Sep 17 00:00:00 2001 From: Effect Bot Date: Fri, 9 Oct 2026 20:29:04 +0000 Subject: [PATCH 5/6] docs(rpc): streamline ping timeout hook comments and changeset --- .changeset/rpc-socket-ping-timeout-hook.md | 2 +- packages/effect/src/rpc/RpcClient.ts | 12 +++--------- 2 files changed, 4 insertions(+), 10 deletions(-) diff --git a/.changeset/rpc-socket-ping-timeout-hook.md b/.changeset/rpc-socket-ping-timeout-hook.md index 78e70559060..18f39363b81 100644 --- a/.changeset/rpc-socket-ping-timeout-hook.md +++ b/.changeset/rpc-socket-ping-timeout-hook.md @@ -2,4 +2,4 @@ "effect": patch --- -Add an `onPingTimeout` option to `RpcClient.makeProtocolSocket` and `RpcClient.layerProtocolSocket`. It runs when a ping timeout drops the connection, before `ConnectionHooks.onDisconnect`, so clients can tell a server that stopped responding from one that closed the connection. +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 5e672bb8cca..f943401e75d 100644 --- a/packages/effect/src/rpc/RpcClient.ts +++ b/packages/effect/src/rpc/RpcClient.ts @@ -1038,8 +1038,6 @@ export const layerProtocolHttp = (options: { * `pingInterval` defaults to 5 seconds. `pingTimeout` defaults to the interval * and measures time since the last decoded server frame, not the last pong. * The timeout is checked on each ping tick; any decoded frame counts as liveness. - * When the timeout drops the connection, `onPingTimeout` runs before - * `ConnectionHooks.onDisconnect`. * * @stability unstable * @category protocols @@ -1080,9 +1078,7 @@ 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 = options?.onPingTimeout ? Effect.ignoreCause(options.onPingTimeout, { @@ -1174,7 +1170,7 @@ export const makeProtocolSocket = (options?: { )) ) }).pipe( - // runs once the socket is closed, so a concurrent socket error cannot cut it short + // Run the hook after socket cleanup so socket errors cannot interrupt it. Effect.tapError((error) => connected && error === pingTimeoutError ? onPingTimeout : Effect.void), Option.isSome(hooks) ? Effect.ensuring(hooks.value.onDisconnect) : identity, Effect.tapCause((cause) => { @@ -1271,9 +1267,7 @@ const makePinger = Effect.fnUntraced(function*(writePing: Effect.Effect * `pingInterval` defaults to 5 seconds. `pingTimeout` defaults to the interval * and measures time since the last decoded server frame, not the last pong. * The timeout is checked on each ping tick; any decoded frame counts as liveness. - * When the timeout drops the connection, `onPingTimeout` runs before - * `ConnectionHooks.onDisconnect`. `retryPolicy` configures retries after socket - * errors. + * `retryPolicy` configures retries after socket errors. * * @stability unstable * @category layers From 21e20ed02e829103af721d972d511b767e3feff7 Mon Sep 17 00:00:00 2001 From: Effect Bot Date: Fri, 9 Oct 2026 20:38:00 +0000 Subject: [PATCH 6/6] refactor(rpc): treat a successful ping race as the timeout --- packages/effect/src/rpc/RpcClient.ts | 38 +++++++++++++--------------- 1 file changed, 17 insertions(+), 21 deletions(-) diff --git a/packages/effect/src/rpc/RpcClient.ts b/packages/effect/src/rpc/RpcClient.ts index f943401e75d..d6d9bb330c2 100644 --- a/packages/effect/src/rpc/RpcClient.ts +++ b/packages/effect/src/rpc/RpcClient.ts @@ -1080,15 +1080,11 @@ export const makeProtocolSocket = (options?: { // Encode each ping with the current connection's parser. const pinger = yield* makePinger(Effect.suspend(() => writer.write(parser.encode(constPing)!)), options) - const onPingTimeout = options?.onPingTimeout ? - Effect.ignoreCause(options.onPingTimeout, { - log: true, - message: "RpcClient onPingTimeout hook failed" - }) : - Effect.void + const onPingTimeout = Effect.ignoreCause(options?.onPingTimeout ?? Effect.void, { + log: true, + message: "RpcClient onPingTimeout hook failed" + }) let currentError: RpcClientError | undefined - let connected = false - let pingTimeoutError: Socket.SocketError | undefined const broadcast = (response: FromServerEncoded) => Effect.forEach(clientIds, (clientId) => writeResponse(clientId, response)) @@ -1141,7 +1137,7 @@ export const makeProtocolSocket = (options?: { yield* Effect.suspend(() => { parser = serialization.makeUnsafe() pinger.reset() - connected = false + let connected = false return Effect.gen(function*() { const { pull } = yield* socket.reader currentError = undefined @@ -1157,21 +1153,21 @@ export const makeProtocolSocket = (options?: { } }).pipe( Effect.scoped, - Effect.raceFirst(Effect.flatMap( - pinger.timeout, - () => - Effect.fail( - pingTimeoutError = 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( - // Run the hook after socket cleanup so socket errors cannot interrupt it. - Effect.tapError((error) => connected && error === pingTimeoutError ? onPingTimeout : Effect.void), Option.isSome(hooks) ? Effect.ensuring(hooks.value.onDisconnect) : identity, Effect.tapCause((cause) => { const error = Cause.findError(cause)