t3-code-android-nightly/patches/effect@4.0.1.patch
Theo Browne 5cebd3b814
fix(clients): dropped connections say why in the client trace (#16200)
Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
2026-10-05 14:31:49 -07:00

139 lines
7.5 KiB
Diff

diff --git a/dist/ai/McpServer.js b/dist/ai/McpServer.js
index 2da594107273980144e84ab9b53636ef34043f6b..9fada51944f58fd0d988ae8f6b18f83e79583a25 100644
--- a/dist/ai/McpServer.js
+++ b/dist/ai/McpServer.js
@@ -1093,7 +1093,7 @@ export const layerHttp = options => {
const methodNotAllowed = request => isAllowedMcpOrigin(request, options.allowedOrigins) ? Effect.succeed(methodNotAllowedResponse) : Effect.succeed(HttpServerResponse.empty({
status: 403
}));
- const routes = Layer.mergeAll(HttpRouter.add("GET", options.path, methodNotAllowed), HttpRouter.add("PUT", options.path, methodNotAllowed), HttpRouter.add("PATCH", options.path, methodNotAllowed), HttpRouter.add("DELETE", options.path, methodNotAllowed), HttpRouter.add("OPTIONS", options.path, methodNotAllowed));
+ const routes = Layer.mergeAll(HttpRouter.add("GET", options.path, methodNotAllowed), HttpRouter.add("PUT", options.path, methodNotAllowed), HttpRouter.add("PATCH", options.path, methodNotAllowed), HttpRouter.add("OPTIONS", options.path, methodNotAllowed));
return Layer.merge(layerWithRuntime(options, "http"), routes).pipe(Layer.provide(layerMcpProtocolHttp(options)), Layer.provide(runtime), Layer.provide(RpcSerialization.layerJsonRpc()));
};
const layerMcpProtocolHttp = options => Layer.effect(RpcServer.Protocol)(Effect.gen(function* () {
@@ -1136,6 +1136,17 @@ const layerMcpProtocolHttp = options => Layer.effect(RpcServer.Protocol)(Effect.
}))) : response;
});
});
+ yield* router.add("DELETE", options.path, request => {
+ if (!isAllowedMcpOrigin(request, options.allowedOrigins)) {
+ return Effect.succeed(HttpServerResponse.empty({
+ status: 403
+ }));
+ }
+ const sessionId = request.headers[MCP_SESSION_ID_HEADER];
+ return Effect.succeed(HttpServerResponse.empty({
+ status: sessionId === undefined ? 400 : runtime.terminateSession(sessionId) ? 204 : 404
+ }));
+ });
return protocol;
}));
const mcpHttpSerialization = /*#__PURE__*/(() => {
diff --git a/dist/ai/internal/mcpRuntime.js b/dist/ai/internal/mcpRuntime.js
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) {
},
effectLogLevel: (clientId, headers, fallback) => stateful?.effectLogLevel(clientId, headers, fallback) ?? fallback,
disconnect: clientId => stateful?.disconnect(clientId),
+ terminateSession: sessionId => stateful?.terminateSession(sessionId) ?? false,
deliveryClientIds: () => stateful?.initializedClientIds() ?? [],
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 9471bdfbcff0643b44df3a08c3701f743a299638..455860316ddfccaf98ff95f9f05da500c331869c 100644
--- a/dist/ai/internal/mcpStatefulRuntime.js
+++ b/dist/ai/internal/mcpStatefulRuntime.js
@@ -46,6 +46,7 @@ export const make = () => {
},
resolve: resolveSession,
resolveSessionId: sessionId => bySessionId.get(sessionId),
+ terminateSession: sessionId => bySessionId.delete(sessionId),
setLogLevel: (level, clientId, headers) => Effect.sync(() => {
const session = resolveSession(clientId, headers);
if (session !== undefined) {
diff --git a/dist/http/HttpClientResponse.js b/dist/http/HttpClientResponse.js
index 42fba673de3b6bdc089c305f2e87ad56c3bfc374..2b4981b8a88606fb93d62a29f21f9556803307c4 100644
--- a/dist/http/HttpClientResponse.js
+++ b/dist/http/HttpClientResponse.js
@@ -195,7 +195,7 @@ class WebHttpClientResponse extends Inspectable.Class {
if (this.cachedCookies) {
return this.cachedCookies;
}
- return this.cachedCookies = Cookies.fromSetCookie(this.source.headers.getSetCookie());
+ return this.cachedCookies = Cookies.fromSetCookie(this.source.headers.getSetCookie?.() ?? []);
}
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 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) {
let recievedPong = true;
+ let missedPongs = 0;
const latch = Latch.makeUnsafe();
const reset = () => {
recievedPong = true;
+ missedPongs = 0;
latch.closeUnsafe();
};
const onPong = () => {
recievedPong = true;
+ missedPongs = 0;
};
yield* Effect.suspend(() => {
- if (!recievedPong) return latch.open;
+ if (!recievedPong) {
+ missedPongs += 1;
+ if (missedPongs >= 3) return latch.open;
+ return writePing;
+ }
recievedPong = false;
+ missedPongs = 0;
return writePing;
}).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 c037015185da1e0c664653b5a496832a9449f3c1..a8b0d8be28bbbd3a0d0d3f1a222c2f074d2d3ad5 100644
--- a/src/http/HttpClientResponse.ts
+++ b/src/http/HttpClientResponse.ts
@@ -341,7 +341,7 @@ class WebHttpClientResponse extends Inspectable.Class implements HttpClientRespo
if (this.cachedCookies) {
return this.cachedCookies
}
- return this.cachedCookies = Cookies.fromSetCookie(this.source.headers.getSetCookie())
+ return this.cachedCookies = Cookies.fromSetCookie(this.source.headers.getSetCookie?.() ?? [])
}
get remoteAddress(): Option.Option<string> {