diff --git a/extensions/telegram/src/api-fetch.live.test.ts b/extensions/telegram/src/api-fetch.live.test.ts new file mode 100644 index 000000000000..ae585e4888f2 --- /dev/null +++ b/extensions/telegram/src/api-fetch.live.test.ts @@ -0,0 +1,73 @@ +import { createServer, type ServerResponse } from "node:http"; +import type { AddressInfo, Socket } from "node:net"; +import { afterEach, describe, expect, it } from "vitest"; +import { fetchTelegramChatId } from "./api-fetch.js"; + +describe("fetchTelegramChatId live HTTP behavior", () => { + const sockets = new Set(); + + afterEach(() => { + for (const socket of sockets) { + socket.destroy(); + } + sockets.clear(); + }); + + it("closes a stalled non-success getChat response body", async () => { + let stalledResponse: ServerResponse | undefined; + let markResponseClosed: () => void = () => undefined; + const responseClosed = new Promise((resolve) => { + markResponseClosed = resolve; + }); + const server = createServer((_req, res) => { + stalledResponse = res; + res.on("close", markResponseClosed); + res.writeHead(503, { "content-type": "text/plain" }); + res.write("service unavailable"); + }); + server.on("connection", (socket) => { + sockets.add(socket); + socket.on("close", () => sockets.delete(socket)); + }); + + await new Promise((resolve) => { + server.listen(0, "127.0.0.1", resolve); + }); + const apiRoot = `http://127.0.0.1:${(server.address() as AddressInfo).port}`; + let captureSettled = false; + const captureFetch: typeof fetch = async (input, init) => { + const response = await fetch(input, init); + const capture = response.clone(); + void capture + .arrayBuffer() + .catch(() => undefined) + .finally(() => { + captureSettled = true; + }); + return response; + }; + + try { + await expect( + fetchTelegramChatId({ token: "abc", chatId: "@user", apiRoot, fetchImpl: captureFetch }), + ).resolves.toBeNull(); + await expect( + Promise.race([ + responseClosed.then(() => "closed"), + new Promise((resolve) => { + setTimeout(() => resolve("stalled"), 1_000); + }), + ]), + ).resolves.toBe("closed"); + await expect.poll(() => captureSettled, { timeout: 1_000 }).toBe(true); + } finally { + stalledResponse?.destroy(); + for (const socket of sockets) { + socket.destroy(); + } + await new Promise((resolve) => { + server.close(() => resolve()); + }); + } + }); +}); diff --git a/extensions/telegram/src/api-fetch.test.ts b/extensions/telegram/src/api-fetch.test.ts index a8a6fc2ec25d..ebe30d93faf2 100644 --- a/extensions/telegram/src/api-fetch.test.ts +++ b/extensions/telegram/src/api-fetch.test.ts @@ -1,7 +1,7 @@ // Telegram tests cover api fetch plugin behavior. import { createRequire } from "node:module"; import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; -import { fetchTelegramChatId } from "./api-fetch.js"; +import { fetchTelegramChatId, lookupTelegramChatId } from "./api-fetch.js"; const TELEGRAM_GETCHAT_JSON_CAP_BYTES = 4 * 1024 * 1024; @@ -177,6 +177,82 @@ describe("fetchTelegramChatId", () => { expect(cancelCount).toBe(1); }); + it("cancels non-success getChat response bodies before returning", async () => { + const cancel = vi.fn(); + let observedSignal: AbortSignal | undefined; + const fetchImpl = vi.fn(async (_url: string | URL | Request, init?: RequestInit) => { + observedSignal = init?.signal ?? undefined; + return new Response( + new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode("service unavailable")); + }, + cancel, + }), + { status: 503 }, + ); + }); + + const result = await Promise.race([ + fetchTelegramChatId({ + token: "abc", + chatId: "@user", + fetchImpl: fetchImpl as unknown as typeof fetch, + }), + new Promise<"stalled">((resolve) => { + setTimeout(() => resolve("stalled"), 250); + }), + ]); + + expect(result).toBeNull(); + expect(cancel).toHaveBeenCalledOnce(); + expect(observedSignal?.aborted).toBe(true); + }); + + it("does not wait for a cloned capture branch before returning", async () => { + let observedSignal: AbortSignal | undefined; + let captureSettled = false; + const fetchImpl = vi.fn(async (_url: string | URL | Request, init?: RequestInit) => { + observedSignal = init?.signal ?? undefined; + const response = new Response( + new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode("service unavailable")); + observedSignal?.addEventListener( + "abort", + () => controller.error(observedSignal?.reason), + { once: true }, + ); + }, + }), + { status: 503 }, + ); + const clone = response.clone(); + void clone + .arrayBuffer() + .catch(() => undefined) + .finally(() => { + captureSettled = true; + }); + return response; + }); + + const result = await Promise.race([ + fetchTelegramChatId({ + token: "abc", + chatId: "@user", + fetchImpl: fetchImpl as unknown as typeof fetch, + }), + new Promise<"stalled">((resolve) => { + setTimeout(() => resolve("stalled"), 250); + }), + ]); + + expect(result).toBeNull(); + await vi.waitFor(() => expect(captureSettled).toBe(true)); + expect(observedSignal?.aborted).toBe(true); + }); + it("keeps the getChat timeout active until the response body read settles", async () => { vi.useFakeTimers(); let observedSignal: AbortSignal | undefined; @@ -218,6 +294,37 @@ describe("fetchTelegramChatId", () => { }); }); +describe("lookupTelegramChatId", () => { + it.each([ + { + name: "success", + setup: () => proxyMocks.undiciFetch.mockResolvedValueOnce(getChatOkResponse(12345)), + expected: "12345", + }, + { + name: "transport failure", + setup: () => proxyMocks.undiciFetch.mockRejectedValueOnce(new Error("network failed")), + expected: null, + }, + ])("closes its owned transport after $name", async ({ setup, expected }) => { + proxyMocks.undiciFetch.mockReset(); + setup(); + + await expect( + lookupTelegramChatId({ + token: "abc", + chatId: "@user", + network: { autoSelectFamily: false }, + }), + ).resolves.toBe(expected); + + const init = proxyMocks.undiciFetch.mock.calls[0]?.[1] as + | (RequestInit & { dispatcher?: { destroyed?: boolean } }) + | undefined; + expect(init?.dispatcher?.destroyed).toBe(true); + }); +}); + describe("undici env proxy semantics", () => { it("uses proxyTls rather than connect for proxied HTTPS transport settings", () => { vi.stubEnv("HTTPS_PROXY", "http://127.0.0.1:7890"); diff --git a/extensions/telegram/src/api-fetch.ts b/extensions/telegram/src/api-fetch.ts index 22f33cab505b..a37ed3fb4450 100644 --- a/extensions/telegram/src/api-fetch.ts +++ b/extensions/telegram/src/api-fetch.ts @@ -2,7 +2,7 @@ import type { TelegramNetworkConfig } from "openclaw/plugin-sdk/config-contracts"; import { buildTimeoutAbortSignal } from "openclaw/plugin-sdk/extension-shared"; import { readResponseWithLimit } from "openclaw/plugin-sdk/response-limit-runtime"; -import { resolveTelegramApiBase, resolveTelegramFetch } from "./fetch.js"; +import { resolveTelegramApiBase, resolveTelegramFetch, resolveTelegramTransport } from "./fetch.js"; import { makeProxyFetch } from "./proxy.js"; import { resolveTelegramRequestTimeoutMs } from "./request-timeouts.js"; @@ -13,6 +13,8 @@ type TelegramGetChatResponse = { result?: { id?: number | string }; }; +// Shipped runtime API. Internal one-shot callers use resolveTelegramTransport +// directly so they can close the dispatcher after the response body settles. export function resolveTelegramChatLookupFetch(params?: { proxyUrl?: string; network?: TelegramNetworkConfig; @@ -31,17 +33,21 @@ export async function lookupTelegramChatId(params: { network?: TelegramNetworkConfig; timeoutSeconds?: unknown; }): Promise { - return fetchTelegramChatId({ - token: params.token, - chatId: params.chatId, - signal: params.signal, - apiRoot: params.apiRoot, - timeoutSeconds: params.timeoutSeconds, - fetchImpl: resolveTelegramChatLookupFetch({ - proxyUrl: params.proxyUrl, - network: params.network, - }), - }); + const proxyUrl = params.proxyUrl?.trim(); + const proxyFetch = proxyUrl ? makeProxyFetch(proxyUrl) : undefined; + const transport = resolveTelegramTransport(proxyFetch, { network: params.network }); + try { + return await fetchTelegramChatId({ + token: params.token, + chatId: params.chatId, + signal: params.signal, + apiRoot: params.apiRoot, + timeoutSeconds: params.timeoutSeconds, + fetchImpl: transport.fetch, + }); + } finally { + await transport.close(); + } } export async function fetchTelegramChatId(params: { @@ -55,8 +61,12 @@ export async function fetchTelegramChatId(params: { const apiBase = resolveTelegramApiBase(params.apiRoot); const url = `${apiBase}/bot${params.token}/getChat?chat_id=${encodeURIComponent(params.chatId)}`; const fetchImpl = params.fetchImpl ?? fetch; + const requestAbortController = new AbortController(); + const requestSignal = params.signal + ? AbortSignal.any([params.signal, requestAbortController.signal]) + : requestAbortController.signal; const timeout = buildTimeoutAbortSignal({ - signal: params.signal, + signal: requestSignal, timeoutMs: resolveTelegramRequestTimeoutMs("getchat", params.timeoutSeconds), operation: "telegram-getchat-lookup", url, @@ -64,6 +74,8 @@ export async function fetchTelegramChatId(params: { try { const res = await fetchImpl(url, timeout.signal ? { signal: timeout.signal } : undefined); if (!res.ok) { + requestAbortController.abort(new Error(`Telegram getChat failed with HTTP ${res.status}`)); + void res.body?.cancel().catch(() => undefined); return null; } let data: TelegramGetChatResponse | null = null;