Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/rpc-socket-liveness.md
Original file line number Diff line number Diff line change
@@ -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.
58 changes: 40 additions & 18 deletions packages/effect/src/rpc/RpcClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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<any, Socket.SocketError> | 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<void>` cannot fail with a typed error or require
* services; defects are logged and ignored so retries can continue.
*/
Expand All @@ -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) =>
Expand All @@ -1079,15 +1088,13 @@ export const makeProtocolSocket = (options?: {
try {
const responses = parser.decode(data) as Array<FromServerEncoded>
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 (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)
Expand Down Expand Up @@ -1202,43 +1209,58 @@ const defaultRetryPolicy = Schedule.min([
Schedule.spaced(5000)
])

const makePinger = Effect.fnUntraced(function*<A, E, R>(writePing: Effect.Effect<A, E, R>) {
let recievedPong = true
const makePinger = Effect.fnUntraced(function*<A, E, R>(writePing: Effect.Effect<A, E, R>, options?: {
readonly pingInterval?: Duration.Input | undefined
readonly pingTimeout?: Duration.Input | undefined
}) {
const clock = yield* Clock.Clock
const interval = options?.pingInterval ?? "5 seconds"
const timeoutMillis = Duration.toMillis(options?.pingTimeout ?? interval)
let lastSeen = clock.currentTimeMillisUnsafe()
const latch = Latch.makeUnsafe()
const onFrame = () => {
lastSeen = clock.currentTimeMillisUnsafe()
}
const reset = () => {
recievedPong = true
onFrame()
latch.closeUnsafe()
}
const onPong = () => {
recievedPong = true
}
yield* Effect.suspend((): Effect.Effect<void, E, R> => {
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<any, Socket.SocketError> | 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<void>` cannot fail with a typed error or require
* services; defects are logged and ignored so retries can continue.
*/
Expand Down
85 changes: 85 additions & 0 deletions packages/effect/test/rpc/RpcClient.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -425,6 +425,91 @@ 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<void>()
const frames = yield* Queue.unbounded<string>()
const events = yield* Queue.unbounded<string>()
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().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<void>()
const frames = yield* Queue.unbounded<string>()
const frameRead = yield* Deferred.make<void>()
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 context = yield* Layer.build(
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))
)
)
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<void>()
Expand Down
Loading