openclaw/.agents/skills/telegram-e2e-userbot/scripts/telegram-test-api-proxy.mjs

439 lines
15 KiB
JavaScript
Executable file

#!/usr/bin/env node
import http from "node:http";
import { Readable } from "node:stream";
// grammY supports client.environment="test", but OpenClaw does not expose it
// through Telegram config. Keep this adapter until that runtime seam exists.
const HOP_BY_HOP_HEADERS = new Set([
"connection",
"keep-alive",
"proxy-authenticate",
"proxy-authorization",
"te",
"trailer",
"transfer-encoding",
"upgrade",
]);
const TELEGRAM_PROXY_ERROR_CODE =
/^(?:ABORT_ERR|EAI_AGAIN|E(?:CONNREFUSED|CONNRESET|HOSTUNREACH|NETUNREACH|NOTFOUND|PIPE|TIMEDOUT)|UND_ERR_(?:ABORTED|BODY_TIMEOUT|CONNECT_TIMEOUT|HEADERS_TIMEOUT|SOCKET))$/u;
export function telegramTestApiPath(pathname) {
const match = pathname.match(/^((?:\/file)?\/bot[^/]+)(\/.*)$/u);
if (!match) throw new Error("Telegram Test Server proxy received an invalid Bot API path.");
return `${match[1]}/test${match[2]}`;
}
function requestHeaders(headers) {
const filtered = {};
for (const [key, value] of Object.entries(headers)) {
if (value !== undefined && key !== "host" && !HOP_BY_HOP_HEADERS.has(key)) {
filtered[key] = value;
}
}
return filtered;
}
function responseHeaders(headers) {
const filtered = {};
for (const [key, value] of headers.entries()) {
if (!HOP_BY_HOP_HEADERS.has(key) && key !== "content-encoding" && key !== "content-length") {
filtered[key] = value;
}
}
return filtered;
}
function telegramApiMethod(pathname) {
return pathname.match(/^\/bot[^/]+\/([^/?]+)/u)?.[1];
}
function reportProxyFailure(error, phase, method) {
let current = error;
let errorCode;
for (let depth = 0; depth < 4 && current && typeof current === "object"; depth += 1) {
if (typeof current.code === "string" && TELEGRAM_PROXY_ERROR_CODE.test(current.code)) {
errorCode = current.code;
break;
}
current = current.cause;
}
// Error messages, URLs, and arbitrary names/codes can contain bot tokens. Emit only
// fixed phases and allowlisted atoms so retained QA logs remain safe to publish.
const safeMethod =
typeof method === "string" && /^[A-Za-z][A-Za-z0-9_]{0,63}$/u.test(method) ? method : "unknown";
console.error(
JSON.stringify({
event: "telegram_test_api_proxy_failure",
phase,
method: safeMethod,
errorClass:
error instanceof TypeError ? "TypeError" : error instanceof Error ? "Error" : "unknown",
errorCode: errorCode ?? "unknown",
}),
);
}
async function drainTelegramTestUpdates(apiRoot, token) {
let offset = 0;
for (;;) {
const response = await fetch(`${apiRoot}/bot${token}/getUpdates`, {
method: "POST",
headers: { "content-type": "application/json" },
body: JSON.stringify({ offset, timeout: 0, allowed_updates: ["message", "edited_message"] }),
signal: AbortSignal.timeout(15_000),
});
const payload = await response.json();
if (!response.ok || payload?.ok !== true || !Array.isArray(payload.result)) {
throw new Error("Telegram Test Bot API getUpdates failed while draining stale updates.");
}
if (payload.result.length === 0) return;
const updateId = payload.result.at(-1)?.update_id;
if (!Number.isSafeInteger(updateId)) {
throw new Error("Telegram Test Bot API getUpdates returned an invalid update.");
}
offset = updateId + 1;
}
}
export async function startTelegramTestApiProxy({
host = "127.0.0.1",
port = 0,
upstream = "https://api.telegram.org",
fetchImpl = fetch,
leaseHealth,
} = {}) {
let responseHold;
let requestRejection;
let heldResponse;
let leaseError;
const upstreamControllers = new Set();
const holdEvents = [];
const rejectionEvents = [];
const requestLog = [];
const methodOrdinals = new Map();
const heldWaiters = new Set();
const sockets = new Set();
const requests = new Set();
let closing;
const assertLeaseHealthy = () => {
if (leaseError) throw leaseError;
leaseHealth?.assertHealthy();
};
const stop = (error) => {
leaseError ??= error;
responseHold = undefined;
requestRejection = undefined;
heldResponse?.release.resolve();
for (const waiter of heldWaiters) waiter.reject(error);
for (const controller of upstreamControllers) controller.abort(error);
for (const socket of sockets) socket.destroy();
};
leaseHealth?.whenUnhealthy.then(stop);
const claimResponseHold = (method, ordinal) => {
if (!responseHold || responseHold.method !== method) return undefined;
if (responseHold.skip > 0) {
responseHold.skip -= 1;
return undefined;
}
const event = { method, ordinal, upstreamAcceptedAt: Date.now(), heldAt: Date.now() };
const release = Promise.withResolvers();
heldResponse = { event, release };
responseHold = undefined;
holdEvents.push(event);
for (const waiter of heldWaiters) {
if (waiter.method === method) waiter.resolve(event);
}
return release.promise;
};
const server = http.createServer((request, response) => {
const work = handleRequest(request, response);
requests.add(work);
work.then(
() => requests.delete(work),
() => requests.delete(work),
);
});
server.on("connection", (socket) => {
sockets.add(socket);
socket.once("close", () => sockets.delete(socket));
});
async function handleRequest(request, response) {
let logged;
let method;
let failurePhase = "request";
const recordDoneAt = () => {
if (logged) {
logged.doneAt ??= Date.now();
}
};
const upstreamController = new AbortController();
upstreamControllers.add(upstreamController);
const abortUpstream = () => upstreamController.abort();
request.once("aborted", abortUpstream);
response.once("close", abortUpstream);
try {
failurePhase = "lease";
assertLeaseHealthy();
failurePhase = "request";
const incoming = new URL(request.url || "/", `http://${host}`);
const upstreamUrl = new URL(upstream);
upstreamUrl.pathname = telegramTestApiPath(incoming.pathname);
upstreamUrl.search = incoming.search;
method = telegramApiMethod(incoming.pathname);
const loggedMethod =
method ?? (incoming.pathname.startsWith("/file/bot") ? "file" : undefined);
const ordinal = (methodOrdinals.get(method) ?? 0) + 1;
// File rows are timing facts only; controls still match the Bot API method.
logged = loggedMethod ? { method: loggedMethod, at: Date.now() } : undefined;
if (logged) {
if (method !== "getUpdates") {
requestLog.push(logged);
}
response.once("finish", recordDoneAt);
response.once("close", recordDoneAt);
}
methodOrdinals.set(method, ordinal);
const hasBody = request.method !== "GET" && request.method !== "HEAD";
let body = hasBody ? request : undefined;
const readBody = async () => {
if (Buffer.isBuffer(body) || !hasBody) return;
const chunks = [];
for await (const chunk of request) chunks.push(Buffer.from(chunk));
body = Buffer.concat(chunks);
};
// Log only the chat kind (private or group), never the chat id, so cross-chat
// flood proof can tell deliveries apart. JSON bodies are small text calls.
if (
logged &&
method !== "getUpdates" &&
String(request.headers["content-type"] ?? "").includes("application/json")
) {
await readBody();
try {
const chatId = Number(JSON.parse(body.toString("utf8")).chat_id);
if (Number.isFinite(chatId) && chatId !== 0) {
logged.chat = chatId < 0 ? "group" : "private";
}
} catch {
// Non-JSON or bodiless calls keep method and timing only.
}
}
const rejection = requestRejection;
if (rejection && rejection.method === method) {
if (rejection.bodyIncludes !== undefined) await readBody();
const matches =
requestRejection === rejection &&
(rejection.bodyIncludes === undefined ||
(Buffer.isBuffer(body) && body.includes(rejection.bodyIncludes)));
if (matches && rejection.skip > 0) {
rejection.skip -= 1;
} else if (matches) {
failurePhase = "lease";
assertLeaseHealthy();
rejection.times -= 1;
if (rejection.times <= 0) requestRejection = undefined;
const flood = rejection.retryAfter !== undefined;
const errorCode = flood ? 429 : 400;
rejectionEvents.push({
method,
ordinal,
rejectedAt: Date.now(),
upstreamForwarded: false,
errorCode,
...(flood ? { retryAfter: rejection.retryAfter } : {}),
});
request.resume();
if (logged) logged.status = errorCode;
response.writeHead(errorCode, { "content-type": "application/json" });
response.end(
JSON.stringify(
flood
? {
ok: false,
error_code: 429,
// retryAfter 0 models a bare 429 without parameters.retry_after.
...(rejection.retryAfter > 0
? {
description: `Too Many Requests: retry after ${rejection.retryAfter}`,
parameters: { retry_after: rejection.retryAfter },
}
: { description: "Too Many Requests" }),
}
: {
ok: false,
error_code: 400,
description: "Bad Request: synthetic Telegram E2E rejection",
},
),
);
return;
}
}
failurePhase = "upstream-fetch";
const result = await fetchImpl(upstreamUrl, {
method: request.method,
headers: requestHeaders(request.headers),
...(hasBody ? { body, duplex: "half" } : {}),
signal: upstreamController.signal,
});
if (logged) logged.status = result.status;
failurePhase = "lease";
assertLeaseHealthy();
failurePhase = "response-stream";
const hold = method ? claimResponseHold(method, ordinal) : undefined;
if (method === "getUpdates") {
try {
const payload = await result.clone().json();
if (Array.isArray(payload?.result) && payload.result.length > 0) {
logged.updates = payload.result.length;
requestLog.push(logged);
}
} catch {
// Preserve malformed responses; only valid update counts are logged.
}
}
if (hold) {
const body = result.body ? Buffer.from(await result.arrayBuffer()) : undefined;
response.writeHead(result.status, responseHeaders(result.headers));
response.write(" ");
const heartbeat = setInterval(() => response.write(" "), 30_000);
heartbeat.unref?.();
response.once("close", () => {
if (heldResponse && !heldResponse.event.releasedAt) {
heldResponse.event.clientClosedAt = Date.now();
}
});
try {
await hold;
} finally {
clearInterval(heartbeat);
}
failurePhase = "lease";
assertLeaseHealthy();
failurePhase = "response-stream";
heldResponse.event.releasedAt = Date.now();
heldResponse = undefined;
response.end(body);
return;
}
failurePhase = "lease";
assertLeaseHealthy();
failurePhase = "response-stream";
response.writeHead(result.status, responseHeaders(result.headers));
if (!result.body) {
response.end();
return;
}
await new Promise((resolve, reject) => {
const readable = Readable.fromWeb(result.body);
const cleanup = () => {
readable.off("error", failed);
response.off("finish", finished);
response.off("close", finished);
};
const failed = (error) => {
cleanup();
reject(error);
};
const finished = () => {
cleanup();
resolve();
};
readable.once("error", failed);
response.once("finish", finished);
response.once("close", finished);
readable.pipe(response);
});
} catch (error) {
recordDoneAt();
reportProxyFailure(error, failurePhase, method);
if (!response.headersSent) {
response.writeHead(502, { "content-type": "application/json" });
}
response.end(
JSON.stringify({ ok: false, description: "Telegram Test Server proxy failed." }),
);
} finally {
request.off("aborted", abortUpstream);
response.off("close", abortUpstream);
upstreamControllers.delete(upstreamController);
}
}
await new Promise((resolve, reject) => {
server.once("error", reject);
server.listen(port, host, () => {
server.off("error", reject);
resolve();
});
});
const address = server.address();
if (!address || typeof address === "string") {
server.close();
throw new Error("Telegram Test Server proxy did not bind a TCP port.");
}
const apiRoot = `http://${host}:${address.port}`;
return {
apiRoot,
drainUpdates: (token) => drainTelegramTestUpdates(apiRoot, token),
rejectNextRequest({ method, skip = 0, bodyIncludes, times = 1, retryAfter }) {
assertLeaseHealthy();
if (requestRejection) throw new Error("A Telegram API request rejection is active.");
requestRejection = { method, skip, bodyIncludes, times, retryAfter };
},
holdNextResponse({ method, skip = 0 }) {
assertLeaseHealthy();
if (responseHold || heldResponse) throw new Error("A Telegram API response hold is active.");
responseHold = { method, skip };
},
waitForHeldResponse(method, timeoutMs) {
assertLeaseHealthy();
if (heldResponse?.event.method === method) return Promise.resolve(heldResponse.event);
return new Promise((resolve, reject) => {
const done = () => {
clearTimeout(timer);
heldWaiters.delete(waiter);
};
const waiter = {
method,
resolve(event) {
done();
resolve(event);
},
reject(error) {
done();
reject(error);
},
};
const timer = setTimeout(
() => waiter.reject(new Error(`Timed out waiting for held Telegram ${method} response.`)),
timeoutMs,
);
heldWaiters.add(waiter);
});
},
releaseHeldResponse() {
assertLeaseHealthy();
if (!heldResponse) throw new Error("No Telegram API response is held.");
const event = heldResponse.event;
heldResponse.release.resolve();
return event;
},
getResponseHoldEvents: () => holdEvents.map((event) => ({ ...event })),
getRequestRejectionEvents: () => rejectionEvents.map((event) => ({ ...event })),
getRequestLog: () => requestLog.map((event) => ({ ...event })),
close: () => {
closing ??= (async () => {
stop(new Error("Telegram Test Server proxy closed."));
await new Promise((resolve) => server.close(() => resolve()));
await Promise.allSettled([...requests]);
})();
return closing;
},
};
}