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-ping-timeout-hook.md
Original file line number Diff line number Diff line change
@@ -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.
60 changes: 41 additions & 19 deletions packages/effect/src/rpc/RpcClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1050,11 +1050,19 @@ export const makeProtocolSocket = (options?: {
readonly retryPolicy?: Schedule.Schedule<any, Socket.SocketError> | 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<void>` 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<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
/**
* 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<void> | undefined
}): Effect.Effect<
Protocol["Service"],
never,
Expand All @@ -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) =>
Expand Down Expand Up @@ -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
}
Expand All @@ -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,
Expand Down Expand Up @@ -1262,11 +1276,19 @@ export const layerProtocolSocket = (options?: {
readonly retryPolicy?: Schedule.Schedule<any, Socket.SocketError> | 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<void>` 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<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
/**
* 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<void> | undefined
}): Layer.Layer<
Protocol,
never,
Expand Down
95 changes: 95 additions & 0 deletions packages/effect/test/rpc/RpcClient.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>()
const events: Array<string> = []
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<string> = []
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<void>()
Expand Down
Loading