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
25 changes: 24 additions & 1 deletion packages/client-runtime/src/connection/supervisor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -802,7 +802,30 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* (
resetRetry: connectedExit.value,
} satisfies AttemptOutcome;
}
return failureFromExit(target, connectedExit, true, connectedForMs >= BACKOFF_RESET_AFTER_MS);
const outcome = failureFromExit(
target,
connectedExit,
true,
connectedForMs >= BACKOFF_RESET_AFTER_MS,
);
if (outcome._tag === "Failure") {
// A live session ending is otherwise invisible in the client trace, so
// record why, and how long it lasted, as its own root span.
yield* Effect.void.pipe(
Effect.withSpan("EnvironmentSupervisor.connectionLost", {
root: true,
attributes: {
"environment.id": target.environmentId,
"environment.label": target.label,
"environment.target.kind": target._tag,
"connection.connected_ms": connectedForMs,
"connection.failure.reason": outcome.failure.error.reason,
"connection.failure.detail": outcome.failure.error.detail,
},
}),
);
}
return outcome;
}, Effect.ensuring(clearLease));

const waitForRetrySignal = Effect.fnUntraced(function* (delayMs: number) {
Expand Down
2 changes: 2 additions & 0 deletions packages/client-runtime/src/errors/transport.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ describe("isTransportConnectionErrorMessage", () => {

it("recognizes connection errors emitted by the Effect RPC session", () => {
expect(isTransportConnectionErrorMessage("Test environment disconnected.")).toBe(true);
expect(isTransportConnectionErrorMessage("Test environment stopped responding.")).toBe(true);
expect(
isTransportConnectionErrorMessage(
"Test environment could not establish a WebSocket connection.",
Expand All @@ -34,6 +35,7 @@ describe("isTransportConnectionErrorMessage", () => {
it("recognizes relay connection errors that carry the network hint", () => {
for (const sentence of [
"Relay environment disconnected.",
"Relay environment stopped responding.",
"Relay environment could not establish a WebSocket connection.",
]) {
expect(isTransportConnectionErrorMessage(`${sentence} ${NETWORK_BLOCKING_HINT}`)).toBe(true);
Expand Down
2 changes: 1 addition & 1 deletion packages/client-runtime/src/errors/transport.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ const TRANSPORT_ERROR_PATTERNS = [
// The RPC session appends the network hint for relay connections. Any other
// trailing text means a different error that the user should still see.
new RegExp(
`\\b(?:is not connected|disconnected|could not establish a WebSocket connection)\\.(?: ${escapeRegExp(NETWORK_BLOCKING_HINT)})?$`,
`\\b(?:is not connected|disconnected|stopped responding|could not establish a WebSocket connection)\\.(?: ${escapeRegExp(NETWORK_BLOCKING_HINT)})?$`,
"i",
),
/\bClientProtocolError\b/i,
Expand Down
5 changes: 4 additions & 1 deletion packages/client-runtime/src/rpc/session.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1108,7 +1108,10 @@ describe("RpcSessionFactory", () => {
yield* TestClock.adjust("5 seconds");
const error = yield* Fiber.join(closedFiber);
expect(error).toBeInstanceOf(ConnectionTransientError);
expect(error).toMatchObject({ reason: "transport" });
expect(error).toMatchObject({
reason: "transport",
detail: "Test environment stopped responding.",
});
}).pipe(Effect.scoped, Effect.provide(TestClock.layer())),
);

Expand Down
16 changes: 11 additions & 5 deletions packages/client-runtime/src/rpc/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -170,18 +170,24 @@ export const make = Effect.fn("RpcSessionFactory.make")(function* (

const connected = yield* Deferred.make<void>();
const disconnected = yield* Deferred.make<never, ConnectionTransientError>();
// Set when the socket closes because pongs stopped, so the failure says so
// instead of looking like the server closed the connection.
const pingTimedOut = yield* Ref.make(false);
const hooks = RpcClient.ConnectionHooks.of({
onConnect: Deferred.succeed(connected, undefined).pipe(Effect.asVoid),
onDisconnect: Deferred.isDone(connected).pipe(
Effect.flatMap((wasConnected) =>
onPingTimeout: Ref.set(pingTimedOut, true),
onDisconnect: Effect.all([Deferred.isDone(connected), Ref.get(pingTimedOut)]).pipe(
Effect.flatMap(([wasConnected, timedOut]) =>
Deferred.fail(
disconnected,
new ConnectionTransientErrorClass({
reason: "transport",
detail: `${
wasConnected
? `${connection.label} disconnected.`
: `${connection.label} could not establish a WebSocket connection.`
!wasConnected
? `${connection.label} could not establish a WebSocket connection.`
: timedOut
? `${connection.label} stopped responding.`
: `${connection.label} disconnected.`
}${networkHint}`,
}),
),
Expand Down
39 changes: 33 additions & 6 deletions patches/effect@4.0.1.patch
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
diff --git a/dist/ai/McpServer.js b/dist/ai/McpServer.js
index 2da5941..9fada51 100644
index 2da594107273980144e84ab9b53636ef34043f6b..9fada51944f58fd0d988ae8f6b18f83e79583a25 100644
--- a/dist/ai/McpServer.js
+++ b/dist/ai/McpServer.js
@@ -1093,7 +1093,7 @@ export const layerHttp = options => {
Expand Down Expand Up @@ -30,7 +30,7 @@ index 2da5941..9fada51 100644
}));
const mcpHttpSerialization = /*#__PURE__*/(() => {
diff --git a/dist/ai/internal/mcpRuntime.js b/dist/ai/internal/mcpRuntime.js
index 06370fd..baa945b 100644
index 06370fd1341ae7a6b681c513a3a3c2866a414b7b..baa945b8b4763b5472c4a900b90b817209caa4bf 100644
--- a/dist/ai/internal/mcpRuntime.js
+++ b/dist/ai/internal/mcpRuntime.js
@@ -336,6 +336,7 @@ export const make = /*#__PURE__*/Effect.fnUntraced(function* (protocols) {
Expand All @@ -42,7 +42,7 @@ index 06370fd..baa945b 100644
canDeliver: (clientId, headers, notification, fallback) => stateful?.canDeliver(clientId, headers, notification, fallback) ?? true,
installHandlers: Effect.fnUntraced(function* (options) {
diff --git a/dist/ai/internal/mcpStatefulRuntime.js b/dist/ai/internal/mcpStatefulRuntime.js
index 9471bdf..4558603 100644
index 9471bdfbcff0643b44df3a08c3701f743a299638..455860316ddfccaf98ff95f9f05da500c331869c 100644
--- a/dist/ai/internal/mcpStatefulRuntime.js
+++ b/dist/ai/internal/mcpStatefulRuntime.js
@@ -46,6 +46,7 @@ export const make = () => {
Expand All @@ -54,7 +54,7 @@ index 9471bdf..4558603 100644
const session = resolveSession(clientId, headers);
if (session !== undefined) {
diff --git a/dist/http/HttpClientResponse.js b/dist/http/HttpClientResponse.js
index 42fba67..2b4981b 100644
index 42fba673de3b6bdc089c305f2e87ad56c3bfc374..2b4981b8a88606fb93d62a29f21f9556803307c4 100644
--- a/dist/http/HttpClientResponse.js
+++ b/dist/http/HttpClientResponse.js
@@ -195,7 +195,7 @@ class WebHttpClientResponse extends Inspectable.Class {
Expand All @@ -66,10 +66,37 @@ index 42fba67..2b4981b 100644
}
get remoteAddress() {
return Option.none();
diff --git a/dist/rpc/RpcClient.d.ts b/dist/rpc/RpcClient.d.ts
index c498ff74510b784dcba791dbf6b8f23dac40a71d..8be9c910a8112a5289159251846128cd27cf32aa 100644
--- a/dist/rpc/RpcClient.d.ts
+++ b/dist/rpc/RpcClient.d.ts
@@ -302,6 +302,8 @@ export declare const layerProtocolWorker: (options: {
declare const ConnectionHooks_base: Context.ServiceClass<ConnectionHooks, "effect/rpc/RpcClient/ConnectionHooks", {
readonly onConnect: Effect.Effect<void>;
readonly onDisconnect: Effect.Effect<void>;
+ /** Runs when the socket is dropped because pongs stopped, before `onDisconnect`. */
+ readonly onPingTimeout?: Effect.Effect<void> | undefined;
}>;
/**
* Represents optional client protocol hooks that run when a transport connects
diff --git a/dist/rpc/RpcClient.js b/dist/rpc/RpcClient.js
index 8af2e2e..65acae4 100644
index 8af2e2efa9311d582add61fef131fd39fd459ab4..723649f05baa61ac33c9bc3634d106caa9641164 100644
--- a/dist/rpc/RpcClient.js
+++ b/dist/rpc/RpcClient.js
@@ -679,11 +679,11 @@ export const makeProtocolSocket = options => Protocol.make(Effect.fnUntraced(fun
yield* processData(frames[i]);
}
}
- }).pipe(Effect.scoped, Effect.raceFirst(Effect.flatMap(pinger.timeout, () => Effect.fail(new Socket.SocketError({
+ }).pipe(Effect.scoped, Effect.raceFirst(Effect.flatMap(pinger.timeout, () => (Option.isSome(hooks) && hooks.value.onPingTimeout ? hooks.value.onPingTimeout : Effect.void).pipe(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, Effect.tapCause(cause => {
const error = Cause.findError(cause);
const hasError = Result.isSuccess(error);
@@ -726,17 +726,25 @@ export const makeProtocolSocket = options => Protocol.make(Effect.fnUntraced(fun
const defaultRetryPolicy = /*#__PURE__*/Schedule.min([/*#__PURE__*/Schedule.exponential(500, 1.5), /*#__PURE__*/Schedule.spaced(5000)]);
const makePinger = /*#__PURE__*/Effect.fnUntraced(function* (writePing) {
Expand Down Expand Up @@ -98,7 +125,7 @@ index 8af2e2e..65acae4 100644
}).pipe(Effect.delay("5 seconds"), Effect.ignore, Effect.forever, Effect.interruptible, Effect.forkScoped);
return {
diff --git a/src/http/HttpClientResponse.ts b/src/http/HttpClientResponse.ts
index c037015..a8b0d8b 100644
index c037015185da1e0c664653b5a496832a9449f3c1..a8b0d8be28bbbd3a0d0d3f1a222c2f074d2d3ad5 100644
--- a/src/http/HttpClientResponse.ts
+++ b/src/http/HttpClientResponse.ts
@@ -341,7 +341,7 @@ class WebHttpClientResponse extends Inspectable.Class implements HttpClientRespo
Expand Down
Loading
Loading