Skip to content

Commit 406a365

Browse files
committed
Observe MCP fetch adapter rejections
1 parent ffaca01 commit 406a365

1 file changed

Lines changed: 42 additions & 34 deletions

File tree

packages/plugins/mcp/src/sdk/connection.ts

Lines changed: 42 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -118,43 +118,51 @@ const abortError = (signal: AbortSignal): unknown => {
118118
return error;
119119
};
120120

121+
const observedPromise = <A>(promise: Promise<A>): Promise<A> => {
122+
void promise.then(undefined, () => undefined);
123+
return promise;
124+
};
125+
121126
const fetchFromHttpClientLayer = (
122127
httpClientLayer: Layer.Layer<HttpClient.HttpClient>,
123128
): FetchLike => {
124-
const execute: FetchLike = async (url, init) => {
125-
const headers = headersFrom(init?.headers);
126-
const requestWithoutBody = HttpClientRequest.make(httpMethodFrom(init?.method))(url, {
127-
headers: recordFromHeaders(headers),
128-
});
129-
const request = await applyBody(requestWithoutBody, headers, init?.body);
130-
const effect = Effect.gen(function* () {
131-
const client = yield* HttpClient.HttpClient;
132-
const response = yield* client.execute(request);
133-
const responseHeaders = new Headers();
134-
for (const [key, value] of Object.entries(response.headers)) {
135-
if (value !== undefined) responseHeaders.set(key, value);
136-
}
137-
const body =
138-
response.status === 204 || response.status === 205 || response.status === 304
139-
? null
140-
: Stream.toReadableStream(response.stream);
141-
return new Response(body, {
142-
status: response.status,
143-
headers: responseHeaders,
144-
});
145-
}).pipe(Effect.provide(httpClientLayer));
146-
const promise = Effect.runPromise(effect);
147-
if (!init?.signal) return promise;
148-
// oxlint-disable-next-line executor/no-promise-reject -- boundary: Fetch-compatible adapter mirrors abort rejection semantics
149-
if (init.signal.aborted) return Promise.reject(abortError(init.signal));
150-
const aborted = new Promise<never>((_, reject) => {
151-
// oxlint-disable-next-line executor/no-promise-reject -- boundary: Fetch-compatible adapter races the Effect request against AbortSignal
152-
init.signal?.addEventListener("abort", () => reject(abortError(init.signal!)), {
153-
once: true,
154-
});
155-
});
156-
return Promise.race([promise, aborted]);
157-
};
129+
const execute: FetchLike = (url, init) =>
130+
observedPromise(
131+
(async () => {
132+
const headers = headersFrom(init?.headers);
133+
const requestWithoutBody = HttpClientRequest.make(httpMethodFrom(init?.method))(url, {
134+
headers: recordFromHeaders(headers),
135+
});
136+
const request = await applyBody(requestWithoutBody, headers, init?.body);
137+
const effect = Effect.gen(function* () {
138+
const client = yield* HttpClient.HttpClient;
139+
const response = yield* client.execute(request);
140+
const responseHeaders = new Headers();
141+
for (const [key, value] of Object.entries(response.headers)) {
142+
if (value !== undefined) responseHeaders.set(key, value);
143+
}
144+
const body =
145+
response.status === 204 || response.status === 205 || response.status === 304
146+
? null
147+
: Stream.toReadableStream(response.stream);
148+
return new Response(body, {
149+
status: response.status,
150+
headers: responseHeaders,
151+
});
152+
}).pipe(Effect.provide(httpClientLayer));
153+
const promise = observedPromise(Effect.runPromise(effect));
154+
if (!init?.signal) return await promise;
155+
// oxlint-disable-next-line executor/no-try-catch-or-throw -- boundary: Fetch-compatible adapter mirrors AbortSignal rejection semantics
156+
if (init.signal.aborted) throw abortError(init.signal);
157+
const aborted = new Promise<never>((_, reject) => {
158+
// oxlint-disable-next-line executor/no-promise-reject -- boundary: Fetch-compatible adapter races the Effect request against AbortSignal
159+
init.signal?.addEventListener("abort", () => reject(abortError(init.signal!)), {
160+
once: true,
161+
});
162+
});
163+
return await Promise.race([promise, aborted]);
164+
})(),
165+
);
158166
return execute;
159167
};
160168

0 commit comments

Comments
 (0)