mirror of
https://github.com/VibedByKaKi/t3-code-android-nightly.git
synced 2026-10-09 03:41:17 +02:00
139 lines
7.5 KiB
Diff
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> {
|