From 3e07f6a2ce9800c0f9927c68238f6be334ee3add Mon Sep 17 00:00:00 2001 From: Tim Smart Date: Mon, 5 Oct 2026 21:20:18 +0000 Subject: [PATCH 1/4] test(rpc): cover socket frame liveness and configurable ping timeout --- packages/effect/test/rpc/RpcClient.test.ts | 89 ++++++++++++++++++++++ 1 file changed, 89 insertions(+) diff --git a/packages/effect/test/rpc/RpcClient.test.ts b/packages/effect/test/rpc/RpcClient.test.ts index 1ca9e1fa1b3..58789878068 100644 --- a/packages/effect/test/rpc/RpcClient.test.ts +++ b/packages/effect/test/rpc/RpcClient.test.ts @@ -425,6 +425,95 @@ describe("RpcClient", () => { assert.strictEqual(error.reason._tag, "SocketReadError") })) + it.effect("keeps in-flight streams alive on non-pong frames without any pongs", () => + Effect.gen(function*() { + const requestSent = yield* Deferred.make() + const frames = yield* Queue.unbounded() + const events = yield* Queue.unbounded() + 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") + }).pipe( + Effect.provideService(Socket.Socket, socket), + 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((event) => Queue.offer(events, event)), + Effect.forkChild + ) + + yield* Deferred.await(requestSent) + for (const event of ["first", "second", "third"]) { + yield* TestClock.adjust("4 seconds") + assert.isUndefined(streamFiber.pollUnsafe()) + yield* Queue.offer(frames, JSON.stringify({ _tag: "Chunk", requestId: "0", values: [event] }) + "\n") + assert.strictEqual(yield* Queue.take(events), event) + } + // Cross another ping tick after the third frame, still without a pong. + yield* TestClock.adjust("4 seconds") + assert.isUndefined(streamFiber.pollUnsafe()) + })) + + it.effect("allows a delayed pong within a custom ping timeout but fails after silence", () => + Effect.gen(function*() { + const requestSent = yield* Deferred.make() + const frames = yield* Queue.unbounded() + const frameRead = yield* Deferred.make() + const write = () => Effect.asVoid(Deferred.succeed(requestSent, void 0)) + const socket = Socket.make({ + reader: Effect.succeed({ + pull: Queue.take(frames).pipe( + Effect.tap(() => Deferred.succeed(frameRead, void 0)), + Effect.map((frame) => [frame] as const) + ), + upgrade: Socket.SocketUpgradeError.unsupported + }), + writer: Effect.succeed({ write, writeAll: write }) + }) + const options = { + pingInterval: "2 seconds", + pingTimeout: "15 seconds", + retryTransientErrors: true, + retryPolicy: Schedule.spaced("1 hour") + } + const context = yield* Layer.build( + RpcClient.layerProtocolSocket(options).pipe( + Layer.provide(RpcSerialization.layerNdjson), + Layer.provide(Layer.succeed(Socket.Socket, socket)) + ) + ) + const client = yield* RpcClient.make(TestGroup).pipe(Effect.provide(context)) + const streamFiber = yield* client.Events().pipe(Stream.runDrain, Effect.exit, Effect.forkChild) + + yield* Deferred.await(requestSent) + // The old fixed timeout fails at 10s, before this delayed pong arrives. + yield* TestClock.adjust("12 seconds") + assert.isUndefined(streamFiber.pollUnsafe()) + yield* Queue.offer(frames, JSON.stringify({ _tag: "Pong" }) + "\n") + yield* Deferred.await(frameRead) + yield* TestClock.adjust("14 seconds") + assert.isUndefined(streamFiber.pollUnsafe()) + + // At 28s the next 2s tick is 16s after the last frame, beyond the 15s timeout. + yield* TestClock.adjust("2 seconds") + assert.isDefined(streamFiber.pollUnsafe()) + const exit = yield* Fiber.join(streamFiber) + assert(Exit.isFailure(exit)) + const error = Cause.squash(exit.cause) + assert.instanceOf(error, RpcClientError) + assert.strictEqual(error.reason._tag, "SocketReadError") + })) + it.effect("fails in-flight streams when transient retries are exhausted", () => Effect.gen(function*() { const requestSent = yield* Deferred.make() From 79bfe70383797e7fc87338bf89273f6f30c21421 Mon Sep 17 00:00:00 2001 From: Tim Smart Date: Mon, 5 Oct 2026 21:24:48 +0000 Subject: [PATCH 2/4] Make RPC socket ping timing configurable and track frame liveness --- .changeset/rpc-socket-liveness.md | 5 +++ packages/effect/src/rpc/RpcClient.ts | 52 ++++++++++++++++++++-------- 2 files changed, 43 insertions(+), 14 deletions(-) create mode 100644 .changeset/rpc-socket-liveness.md diff --git a/.changeset/rpc-socket-liveness.md b/.changeset/rpc-socket-liveness.md new file mode 100644 index 00000000000..efb977b5f31 --- /dev/null +++ b/.changeset/rpc-socket-liveness.md @@ -0,0 +1,5 @@ +--- +"effect": patch +--- + +Add configurable `pingInterval` and `pingTimeout` to RPC socket clients, counting any decoded server frame as liveness. Forward `retryPolicy` through `layerProtocolSocket`. The default ping interval and timeout remain 5 seconds. diff --git a/packages/effect/src/rpc/RpcClient.ts b/packages/effect/src/rpc/RpcClient.ts index 3272af43e93..210bec6872f 100644 --- a/packages/effect/src/rpc/RpcClient.ts +++ b/packages/effect/src/rpc/RpcClient.ts @@ -13,8 +13,9 @@ */ import type { NonEmptyReadonlyArray } from "../Array.ts" import * as Cause from "../Cause.ts" +import * as Clock from "../Clock.ts" import * as Context from "../Context.ts" -import type * as Duration from "../Duration.ts" +import * as Duration from "../Duration.ts" import * as Effect from "../Effect.ts" import * as Exit from "../Exit.ts" import * as Fiber from "../Fiber.ts" @@ -1030,16 +1031,24 @@ export const layerProtocolHttp = (options: { * `RpcSerialization`, connection hooks, ping timeouts, and the configured retry * policy. * + * **Details** + * + * `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. + * * @stability unstable * @category protocols * @since 4.0.0 */ export const makeProtocolSocket = (options?: { + readonly pingInterval?: Duration.Input | undefined + readonly pingTimeout?: Duration.Input | undefined readonly retryTransientErrors?: boolean | undefined readonly retryPolicy?: Schedule.Schedule | undefined /** * Runs for each retried `SocketOpenError` when `retryTransientErrors` is enabled. - * A missed pong fails in-flight calls and is not reported through this hook. + * 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. */ @@ -1062,7 +1071,7 @@ export const makeProtocolSocket = (options?: { // `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. - const pinger = yield* makePinger(Effect.suspend(() => writer.write(parser.encode(constPing)!))) + const pinger = yield* makePinger(Effect.suspend(() => writer.write(parser.encode(constPing)!)), options) let currentError: RpcClientError | undefined const broadcast = (response: FromServerEncoded) => @@ -1079,13 +1088,13 @@ export const makeProtocolSocket = (options?: { try { const responses = parser.decode(data) as Array if (responses.length === 0) return Effect.void + pinger.onFrame() let i = 0 return Effect.whileLoop({ while: () => i < responses.length, body: () => { const response = responses[i++] if (response._tag === "Pong") { - pinger.onPong() return Effect.void } if (Object.hasOwn(response, "requestId")) { @@ -1202,43 +1211,58 @@ const defaultRetryPolicy = Schedule.min([ Schedule.spaced(5000) ]) -const makePinger = Effect.fnUntraced(function*(writePing: Effect.Effect) { - let recievedPong = true +const makePinger = Effect.fnUntraced(function*(writePing: Effect.Effect, options?: { + readonly pingInterval?: Duration.Input | undefined + readonly pingTimeout?: Duration.Input | undefined +}) { + const clock = yield* Clock.Clock + const interval = Duration.fromInputUnsafe(options?.pingInterval ?? "5 seconds") + const timeoutMillis = Duration.toMillis(Duration.fromInputUnsafe(options?.pingTimeout ?? interval)) + let lastSeen = clock.currentTimeMillisUnsafe() const latch = Latch.makeUnsafe() const reset = () => { - recievedPong = true + lastSeen = clock.currentTimeMillisUnsafe() latch.closeUnsafe() } - const onPong = () => { - recievedPong = true + const onFrame = () => { + lastSeen = clock.currentTimeMillisUnsafe() } yield* Effect.suspend((): Effect.Effect => { - if (!recievedPong) return latch.open - recievedPong = false + if (clock.currentTimeMillisUnsafe() - lastSeen > timeoutMillis) return latch.open return writePing }).pipe( - Effect.delay("5 seconds"), + Effect.delay(interval), Effect.ignore, Effect.forever, Effect.interruptible, Effect.forkScoped ) - return { timeout: latch.await, reset, onPong } as const + return { timeout: latch.await, reset, onFrame } as const }) /** * Provides a client `Protocol` backed by the current `Socket` and * `RpcSerialization` services. * + * **Details** + * + * `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. + * * @stability unstable * @category layers * @since 4.0.0 */ export const layerProtocolSocket = (options?: { + readonly pingInterval?: Duration.Input | undefined + readonly pingTimeout?: Duration.Input | undefined readonly retryTransientErrors?: boolean | undefined + readonly retryPolicy?: Schedule.Schedule | undefined /** * Runs for each retried `SocketOpenError` when `retryTransientErrors` is enabled. - * A missed pong fails in-flight calls and is not reported through this hook. + * 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. */ From 362943e7a032a4991f81ca490e873fecd6b9ed59 Mon Sep 17 00:00:00 2001 From: Tim Smart Date: Mon, 5 Oct 2026 21:27:18 +0000 Subject: [PATCH 3/4] Preserve RPC ping duration literals in test options --- packages/effect/test/rpc/RpcClient.test.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/effect/test/rpc/RpcClient.test.ts b/packages/effect/test/rpc/RpcClient.test.ts index 58789878068..c82c20b022c 100644 --- a/packages/effect/test/rpc/RpcClient.test.ts +++ b/packages/effect/test/rpc/RpcClient.test.ts @@ -485,7 +485,7 @@ describe("RpcClient", () => { pingTimeout: "15 seconds", retryTransientErrors: true, retryPolicy: Schedule.spaced("1 hour") - } + } as const const context = yield* Layer.build( RpcClient.layerProtocolSocket(options).pipe( Layer.provide(RpcSerialization.layerNdjson), From 3bfc94925d9dea9b13efb4889b893753d8594a80 Mon Sep 17 00:00:00 2001 From: Tim Smart Date: Mon, 5 Oct 2026 21:36:06 +0000 Subject: [PATCH 4/4] Simplify RPC pinger duration handling and liveness tests --- packages/effect/src/rpc/RpcClient.ts | 16 +++++++--------- packages/effect/test/rpc/RpcClient.test.ts | 16 ++++++---------- 2 files changed, 13 insertions(+), 19 deletions(-) diff --git a/packages/effect/src/rpc/RpcClient.ts b/packages/effect/src/rpc/RpcClient.ts index 210bec6872f..93a1e1ac849 100644 --- a/packages/effect/src/rpc/RpcClient.ts +++ b/packages/effect/src/rpc/RpcClient.ts @@ -1094,9 +1094,7 @@ export const makeProtocolSocket = (options?: { while: () => i < responses.length, body: () => { const response = responses[i++] - if (response._tag === "Pong") { - return Effect.void - } + if (response._tag === "Pong") return Effect.void if (Object.hasOwn(response, "requestId")) { const requestId = (response as FromServerEncoded & { readonly requestId: string | number }).requestId const clientId = requestClientMap.get(requestId) @@ -1216,17 +1214,17 @@ const makePinger = Effect.fnUntraced(function*(writePing: Effect.Effect readonly pingTimeout?: Duration.Input | undefined }) { const clock = yield* Clock.Clock - const interval = Duration.fromInputUnsafe(options?.pingInterval ?? "5 seconds") - const timeoutMillis = Duration.toMillis(Duration.fromInputUnsafe(options?.pingTimeout ?? interval)) + const interval = options?.pingInterval ?? "5 seconds" + const timeoutMillis = Duration.toMillis(options?.pingTimeout ?? interval) let lastSeen = clock.currentTimeMillisUnsafe() const latch = Latch.makeUnsafe() - const reset = () => { - lastSeen = clock.currentTimeMillisUnsafe() - latch.closeUnsafe() - } const onFrame = () => { lastSeen = clock.currentTimeMillisUnsafe() } + const reset = () => { + onFrame() + latch.closeUnsafe() + } yield* Effect.suspend((): Effect.Effect => { if (clock.currentTimeMillisUnsafe() - lastSeen > timeoutMillis) return latch.open return writePing diff --git a/packages/effect/test/rpc/RpcClient.test.ts b/packages/effect/test/rpc/RpcClient.test.ts index c82c20b022c..2a930062299 100644 --- a/packages/effect/test/rpc/RpcClient.test.ts +++ b/packages/effect/test/rpc/RpcClient.test.ts @@ -438,9 +438,7 @@ describe("RpcClient", () => { }), writer: Effect.succeed({ write, writeAll: write }) }) - const protocol = yield* RpcClient.makeProtocolSocket({ - retryPolicy: Schedule.spaced("1 hour") - }).pipe( + const protocol = yield* RpcClient.makeProtocolSocket().pipe( Effect.provideService(Socket.Socket, socket), Effect.provide(RpcSerialization.layerNdjson) ) @@ -480,14 +478,12 @@ describe("RpcClient", () => { }), writer: Effect.succeed({ write, writeAll: write }) }) - const options = { - pingInterval: "2 seconds", - pingTimeout: "15 seconds", - retryTransientErrors: true, - retryPolicy: Schedule.spaced("1 hour") - } as const const context = yield* Layer.build( - RpcClient.layerProtocolSocket(options).pipe( + RpcClient.layerProtocolSocket({ + pingInterval: "2 seconds", + pingTimeout: "15 seconds", + retryPolicy: Schedule.spaced("1 hour") + }).pipe( Layer.provide(RpcSerialization.layerNdjson), Layer.provide(Layer.succeed(Socket.Socket, socket)) )