feat: add Telegram push support with multi-target routes

Add a platform-aware route target system so a single route can forward
to several destinations at once (e.g. a Discord channel and a Telegram
group). Route.target becomes Route.targets[] with per-entry platform,
channelId/threadId for Discord and chatId/topicId for Telegram; the
legacy single-target format is normalized on load and accepted by the
admin API.

Implement the Telegram driver with HTML rendering and Bot API
sendMessage (chat_id + message_thread_id for topics, retry on 429/5xx),
plus /gh commands served over POST /telegram/webhook: login, logout,
comment, merge and close. The comment/merge/close commands resolve the
issue or PR from the replied-to notification message. OAuth binding now
stores a D1 telegram_links mapping and replies with a confirmation.

Sync the Telegram webhook from the scheduled trigger via setWebhook.
This commit is contained in:
RhenCloud 2026-08-03 05:17:59 +08:00
parent dcbe93be91
commit bd7a8f2632
No known key found for this signature in database
GPG key ID: A574A617378C4E0B
39 changed files with 1165 additions and 162 deletions

View file

@ -52,7 +52,7 @@ const sampleRoutes: Route[] = [
name: "Backend PRs",
enabled: true,
filters: [{ type: "event", match: "pull_request" }],
target: { channelId: "111" },
targets: [{ channelId: "111" }],
},
];
@ -133,6 +133,6 @@ describe("config routes persistence", () => {
const second = await loadConfig(env);
expect(second.routes).toHaveLength(1);
expect(second.routes[0]!.id).toBe("backend-prs");
expect(second.routes[0]!.target.channelId).toBe("111");
expect(second.routes[0]!.targets[0]!.channelId).toBe("111");
});
});

View file

@ -129,7 +129,7 @@ describe("dispatchEvent fallback routing", () => {
name: "Regular Push",
enabled: true,
filters: [{ type: "event", match: "push" }],
target: { channelId: "111" },
targets: [{ channelId: "111" }],
},
{
id: "catch-all",
@ -137,7 +137,7 @@ describe("dispatchEvent fallback routing", () => {
enabled: true,
filters: [],
fallback: true,
target: { channelId: "222" },
targets: [{ channelId: "222" }],
},
];

View file

@ -7,7 +7,7 @@ const route: Route = {
name: "Test",
enabled: true,
filters: [],
target: { channelId: "111" },
targets: [{ channelId: "111" }],
};
function event(ev: string, payload: Record<string, unknown>): WebhookEvent {

View file

@ -0,0 +1,201 @@
import { describe, it, expect, afterEach } from "bun:test";
import { handleTelegramUpdate } from "../drivers/telegram/commands";
import { handleTelegramWebhookRequest } from "../drivers/telegram/updates";
import type { Env } from "../types";
function mockFetch(handler: (url: string, init?: RequestInit) => Response): void {
globalThis.fetch = (input: RequestInfo | URL, init?: RequestInit): Promise<Response> =>
Promise.resolve(handler(String(input), init));
}
const restoredFetch = globalThis.fetch;
afterEach(() => {
globalThis.fetch = restoredFetch;
});
function createMockKV(): KVNamespace {
const store = new Map<string, string>();
return {
get: async (key: string, type?: string) => {
const v = store.get(key);
if (v === undefined) return null;
return type === "json" ? JSON.parse(v) : v;
},
put: async (key: string, value: string) => {
store.set(key, value);
},
delete: async (key: string) => {
store.delete(key);
},
list: async () => ({ keys: [] }),
} as unknown as KVNamespace;
}
function createMockDB(): D1Database {
const links = new Map<string, string>();
return {
prepare: (sql: string) => ({
bind: (...args: unknown[]) => ({
run: async (): Promise<{ success: boolean }> => {
const m = sql.match(/INSERT OR REPLACE INTO telegram_links \(telegram_user_id, github_user_id\) VALUES \(\?, \?\)/);
if (m) links.set(String(args[0]), String(args[1]));
const del = sql.match(/DELETE FROM telegram_links WHERE telegram_user_id = \?/);
if (del) links.delete(String(args[0]));
return { success: true };
},
all: async (): Promise<{ results: Array<Record<string, unknown>> }> => {
const sel = sql.match(/SELECT github_user_id FROM telegram_links WHERE telegram_user_id = \?/);
if (sel) {
const val = links.get(String(args[0]));
return { results: val ? [{ github_user_id: val }] : [] };
}
return { results: [] };
},
first: async () => null,
}),
}),
} as unknown as D1Database;
}
function createEnv(): Env {
return {
GITHUB_WEBHOOK_SECRET: "secret",
GITHUB_CLIENT_ID: "client-id",
TELEGRAM_TOKEN: "tg-token",
TELEGRAM_WEBHOOK_SECRET: "wh-secret",
KV: createMockKV(),
DB: createMockDB(),
} as Env;
}
function reply(chatId: string, topicId?: number): Record<string, unknown> {
return {
message_id: 1,
from: { id: 111, first_name: "Rhen" },
chat: { id: chatId, type: "supergroup" },
message_thread_id: topicId,
};
}
describe("telegram-commands /gh login", () => {
it("stores state with telegramUserId and replies with OAuth URL", async () => {
let sentBody: Record<string, unknown> | undefined;
mockFetch((_url, init) => {
sentBody = JSON.parse(String(init!.body)) as Record<string, unknown>;
return new Response(JSON.stringify({ ok: true, result: { message_id: 1 } }), { status: 200 });
});
const env = createEnv();
await handleTelegramUpdate(env, {
message: { ...reply("-100123"), text: "/gh login" },
});
expect(sentBody?.chat_id).toBe("-100123");
expect(String(sentBody?.text)).toContain("github.com/login/oauth/authorize");
});
});
describe("telegram-commands /gh logout", () => {
it("replies bound/unbound message", async () => {
let sentText = "";
mockFetch((_url, init) => {
sentText = String((JSON.parse(String(init!.body)) as Record<string, unknown>).text);
return new Response(JSON.stringify({ ok: true, result: { message_id: 1 } }), { status: 200 });
});
const env = createEnv();
await handleTelegramUpdate(env, {
message: { ...reply("-100123"), text: "/gh logout" },
});
expect(sentText).toContain("已解绑");
});
});
describe("telegram-commands /gh comment", () => {
it("replies when no reply_to_message link present", async () => {
let sentText = "";
mockFetch((_url, init) => {
sentText = String((JSON.parse(String(init!.body)) as Record<string, unknown>).text);
return new Response(JSON.stringify({ ok: true, result: { message_id: 1 } }), { status: 200 });
});
const env = createEnv();
await handleTelegramUpdate(env, {
message: { ...reply("-100123"), text: "/gh comment hello" },
});
expect(sentText).toContain("还没有绑定");
});
it("parses the replied-to GitHub link as target", async () => {
const env = createEnv();
const { saveTelegramLink } = await import("../github/store");
await saveTelegramLink(env.DB, "111", "111980217");
await env.KV.put(
"token:111980217",
JSON.stringify({
userId: "111980217",
accessToken: "ghu_test",
expiresAt: Date.now() + 3600_000,
}),
);
const calls: Array<{ url: string; body: string }> = [];
mockFetch((url, init) => {
calls.push({ url, body: String(init!.body) });
if (String(url).endsWith("/sendMessage")) {
return new Response(JSON.stringify({ ok: true, result: { message_id: 1 } }), { status: 200 });
}
return new Response(
JSON.stringify({ html_url: "https://github.com/acme/widget/issues/7#issuecomment-9" }),
{ status: 200 },
);
});
await handleTelegramUpdate(env, {
message: {
...reply("-100123"),
text: "/gh comment hello",
reply_to_message: {
...reply("-100123"),
text: "acme/widget#7: Add feature",
entities: [{ type: "text_link", url: "https://github.com/acme/widget/issues/7" }],
},
},
});
const ghCall = calls.find((c) => c.url.includes("/repos/"));
expect(ghCall).toBeDefined();
expect(ghCall!.url).toContain("/repos/acme/widget/issues/7/comments");
const ghBody = JSON.parse(ghCall!.body) as Record<string, unknown>;
expect(ghBody.body).toBe("hello");
});
});
describe("telegram-updates webhook", () => {
it("rejects requests without the secret token", async () => {
const env = createEnv();
const res = await handleTelegramWebhookRequest(
new Request("https://example.com/telegram/webhook", { method: "POST", body: "{}" }),
env,
);
expect(res.status).toBe(401);
});
it("accepts requests with the correct secret token", async () => {
mockFetch(() => new Response(JSON.stringify({ ok: true, result: { message_id: 1 } }), { status: 200 }));
const env = createEnv();
const res = await handleTelegramWebhookRequest(
new Request("https://example.com/telegram/webhook", {
method: "POST",
headers: { "X-Telegram-Bot-Api-Secret-Token": "wh-secret" },
body: JSON.stringify({ update_id: 1 }),
}),
env,
);
expect(res.status).toBe(200);
expect(await res.text()).toBe("ok");
});
});

View file

@ -0,0 +1,92 @@
import { describe, it, expect, afterEach } from "bun:test";
import { sendMessage } from "../drivers/telegram/rest";
import { renderNeutralMessage } from "../drivers/telegram/render";
import type { NeutralMessage } from "../types";
function mockFetch(handler: (url: string, init?: RequestInit) => Response): void {
globalThis.fetch = (input: RequestInfo | URL, init?: RequestInit): Promise<Response> =>
Promise.resolve(handler(String(input), init));
}
const restoredFetch = globalThis.fetch;
afterEach(() => {
globalThis.fetch = restoredFetch;
});
describe("telegram renderNeutralMessage", () => {
it("renders title, fields and footer as HTML", () => {
const message: NeutralMessage = {
title: "acme/widget: Add feature",
url: "https://github.com/acme/widget",
fields: [{ name: "Status", value: "success" }],
footer: "acme/widget",
};
const out = renderNeutralMessage(message);
expect(out).toContain('<b><a href="https://github.com/acme/widget">acme/widget: Add feature</a></b>');
expect(out).toContain("<b>Status</b>: success");
expect(out).toContain("<i>acme/widget</i>");
});
it("escapes HTML special characters", () => {
const out = renderNeutralMessage({
title: "a <b> & \"c\"",
fields: [{ name: "body", value: "<script>alert(1)</script>" }],
});
expect(out).not.toContain("<b>acme");
expect(out).toContain("&lt;script&gt;alert(1)&lt;/script&gt;");
});
});
describe("telegram-rest sendMessage", () => {
it("posts to the bot API with chat_id and parse_mode", async () => {
let capturedUrl = "";
let capturedInit: RequestInit | undefined;
mockFetch((url, init) => {
capturedUrl = url;
capturedInit = init;
return new Response(JSON.stringify({ ok: true, result: { message_id: 42 } }), {
status: 200,
});
});
const result = await sendMessage("token-abc", "-100123", "hello");
expect(result.ok).toBe(true);
expect(result.messageId).toBe("42");
expect(capturedUrl).toBe("https://api.telegram.org/bottoken-abc/sendMessage");
const body = JSON.parse(String(capturedInit!.body)) as Record<string, unknown>;
expect(body.chat_id).toBe("-100123");
expect(body.parse_mode).toBe("HTML");
expect(body.message_thread_id).toBeUndefined();
});
it("includes message_thread_id when topicId is given", async () => {
let capturedInit: RequestInit | undefined;
mockFetch((_url, init) => {
capturedInit = init;
return new Response(JSON.stringify({ ok: true, result: { message_id: 1 } }), {
status: 200,
});
});
await sendMessage("t", "-100123", "hello", "999");
const body = JSON.parse(String(capturedInit!.body)) as Record<string, unknown>;
expect(body.message_thread_id).toBe(999);
});
it("returns error on non-ok response", async () => {
mockFetch(() =>
new Response(JSON.stringify({ ok: false, description: "chat not found" }), { status: 400 }),
);
const result = await sendMessage("t", "-100123", "hello");
expect(result.ok).toBe(false);
expect(result.error).toContain("chat not found");
});
it("returns error when token is missing", async () => {
const result = await sendMessage("", "-100123", "hello");
expect(result.ok).toBe(false);
expect(result.errorCode).toBe("NO_TOKEN");
});
});

View file

@ -51,7 +51,7 @@ describe("matchRoute", () => {
name: "Test",
enabled: true,
filters: [],
target: { channelId: "123" },
targets: [{ channelId: "123" }],
};
it("matches all events when no filters", () => {

View file

@ -8,13 +8,26 @@ let configCache: { config: Config; expiresAt: number } | null = null;
export async function loadRoutes(kv: KVNamespace): Promise<Route[]> {
try {
const stored = await kv.get<Route[]>(ROUTES_KEY, "json");
if (stored) return stored;
if (stored) return normalizeRoutes(stored);
} catch (err) {
log.warn({ err }, "Failed to load routes from KV");
}
return [];
}
/**
* Migrate legacy single-target routes (`target`) to the array form (`targets`).
*/
function normalizeRoutes(routes: Route[]): Route[] {
return routes.map((r) => {
if (r.targets && r.targets.length > 0) return r;
const legacy = (r as Route & { target?: Route["targets"][number] }).target;
if (!legacy) return r;
const { target: _target, ...rest } = r as Route & { target?: Route["targets"][number] };
return { ...rest, targets: [legacy] };
});
}
export async function saveRoutes(kv: KVNamespace, routes: Route[]): Promise<void> {
await kv.put(ROUTES_KEY, JSON.stringify(routes));
configCache = null;

View file

@ -37,59 +37,70 @@ export async function dispatchEvent(config: Config, event: WebhookEvent, env: En
return matchRoute(route, event);
})
.map(async (route) => {
const target = route.target.threadId
? `${route.target.channelId}/${route.target.threadId}`
: route.target.channelId;
const targets = route.targets && route.targets.length > 0 ? route.targets : [];
if (targets.length === 0) return;
const base: {
ts: number;
routeId: string;
event: string;
repo: string | undefined;
target: string;
deliveryId: string | undefined;
actor: string | undefined;
action: string | undefined;
} = {
ts: Date.now(),
routeId: route.id,
event: event.event,
repo: (event.payload.repository as { full_name?: string } | undefined)?.full_name,
target,
deliveryId: event.deliveryId,
actor: (event.payload.sender as { login?: string } | undefined)?.login,
action: (event.payload.action as string | undefined),
};
const tr = trMap.get(route.lang ?? "en")!;
const group = route.groupId ? groupById.get(route.groupId) : undefined;
const showEmoji = group?.emoji !== false;
const message = formatEvent(route, event, tr, showEmoji);
const started = Date.now();
try {
const tr = trMap.get(route.lang ?? "en")!;
const group = route.groupId ? groupById.get(route.groupId) : undefined;
const showEmoji = group?.emoji !== false;
const message = formatEvent(route, event, tr, showEmoji);
const driver = getDriver(route.target);
const result = await driver.send(message, route.target, env);
const durationMs = Date.now() - started;
if (!result.ok) throw new Error(result.error ?? "Send failed");
await recordSend(env.DB, {
...base,
ok: true,
status: result.status,
messageId: result.messageId,
platform: driver.id,
attempts: result.attempts,
durationMs,
errorCode: result.errorCode,
});
} catch (err) {
const durationMs = Date.now() - started;
await recordSend(env.DB, {
...base,
ok: false,
error: err instanceof Error ? err.message : String(err),
durationMs,
});
log.error({ routeId: route.id, err }, "Route failed");
for (const target of targets) {
const targetStr =
target.platform === "telegram"
? target.topicId
? `${target.chatId}/${target.topicId}`
: (target.chatId ?? "")
: target.threadId
? `${target.channelId}/${target.threadId}`
: (target.channelId ?? "");
const base: {
ts: number;
routeId: string;
event: string;
repo: string | undefined;
target: string;
deliveryId: string | undefined;
actor: string | undefined;
action: string | undefined;
} = {
ts: Date.now(),
routeId: route.id,
event: event.event,
repo: (event.payload.repository as { full_name?: string } | undefined)?.full_name,
target: targetStr,
deliveryId: event.deliveryId,
actor: (event.payload.sender as { login?: string } | undefined)?.login,
action: (event.payload.action as string | undefined),
};
const started = Date.now();
try {
const driver = getDriver(target);
const result = await driver.send(message, target, env);
const durationMs = Date.now() - started;
if (!result.ok) throw new Error(result.error ?? "Send failed");
await recordSend(env.DB, {
...base,
ok: true,
status: result.status,
messageId: result.messageId,
platform: driver.id,
attempts: result.attempts,
durationMs,
errorCode: result.errorCode,
});
} catch (err) {
const durationMs = Date.now() - started;
await recordSend(env.DB, {
...base,
ok: false,
error: err instanceof Error ? err.message : String(err),
durationMs,
});
log.error({ routeId: route.id, target: targetStr, err }, "Route failed");
}
}
});

View file

@ -1,4 +1,4 @@
import type { Route, Env, NeutralMessage } from "../../types";
import type { RouteTarget, Env, NeutralMessage } from "../../types";
import type { PlatformDriver, SendResult } from "../types";
import { sendMessage } from "./rest";
import { renderNeutralMessage } from "./render";
@ -6,8 +6,12 @@ import { renderNeutralMessage } from "./render";
export class DiscordDriver implements PlatformDriver {
readonly id = "discord";
async send(message: NeutralMessage, target: Route["target"], env: Env): Promise<SendResult> {
async send(message: NeutralMessage, target: RouteTarget, env: Env): Promise<SendResult> {
const channelId = target.channelId ?? "";
if (!channelId) {
return { ok: false, error: "target.channelId is required", errorCode: "NO_TARGET" };
}
const token = env.DISCORD_TOKEN ?? "";
return sendMessage(token, target.channelId, renderNeutralMessage(message), target.threadId);
return sendMessage(token, channelId, renderNeutralMessage(message), target.threadId);
}
}

View file

@ -1,4 +1,4 @@
import type { Route } from "../types";
import type { RouteTarget } from "../types";
import type { PlatformDriver } from "./types";
import { DiscordDriver } from "./discord";
import { TelegramDriver } from "./telegram";
@ -8,8 +8,8 @@ const drivers: Record<string, PlatformDriver> = {
telegram: new TelegramDriver(),
};
export function getDriver(target: Route["target"]): PlatformDriver {
const platform = (target as { platform?: string }).platform ?? "discord";
export function getDriver(target: RouteTarget): PlatformDriver {
const platform = target.platform ?? "discord";
const driver = drivers[platform];
if (!driver) throw new Error(`No driver for platform "${platform}"`);
return driver;

View file

@ -0,0 +1,231 @@
import { log } from "../../lib/log";
import {
getOAuthURL,
commentAsUser,
mergePullRequestAsUser,
closePullRequestAsUser,
} from "../../github/oauth";
import { getTelegramLink, removeTelegramLink } from "../../github/store";
import type { Env } from "../../types";
import { sendMessage } from "./rest";
const GITHUB_ISSUE_RE = /github\.com\/([^/\s]+)\/([^/\s]+)\/(?:issues|pull)\/(\d+)/;
const GITHUB_PR_RE = /github\.com\/([^/\s]+)\/([^/\s]+)\/pull\/(\d+)/;
interface TelegramMessage {
message_id?: number;
text?: string;
from?: { id?: number; first_name?: string; username?: string };
chat?: { id?: number; type?: string; title?: string };
date?: number;
message_thread_id?: number;
entities?: Array<{ type?: string; url?: string; offset?: number; length?: number }>;
reply_to_message?: TelegramMessage;
}
interface Target {
owner: string;
repo: string;
number: number;
}
function chatIdOf(msg: TelegramMessage): string | null {
return msg.chat?.id != null ? String(msg.chat.id) : null;
}
function userIdOf(msg: TelegramMessage): string | null {
return msg.from?.id != null ? String(msg.from.id) : null;
}
/** Extract a GitHub issue/PR link from a message (entities text_link or raw text). */
function extractTarget(msg: TelegramMessage, prOnly = false): Target | null {
const urls: string[] = [];
for (const ent of msg.entities ?? []) {
if (ent.type === "text_link" && ent.url) urls.push(ent.url);
}
if (msg.text) {
for (const u of msg.text.match(/https?:\/\/github\.com\/[^\s]+/g) ?? []) urls.push(u);
}
for (const url of urls) {
const re = prOnly ? GITHUB_PR_RE : GITHUB_ISSUE_RE;
const m = url.match(re);
if (m) return { owner: m[1], repo: m[2], number: Number(m[3]) };
}
return null;
}
async function reply(env: Env, chatId: string, topicId: string | undefined, text: string): Promise<void> {
await sendMessage(env.TELEGRAM_TOKEN ?? "", chatId, text, topicId);
}
function errText(err: unknown): string {
const t = err instanceof Error ? err.message : String(err);
if (t === "GITHUB_TOKEN_EXPIRED") return "GitHub 授权已过期或无效,请重新使用 /gh login 绑定。";
if (t === "GITHUB_FORBIDDEN") return "GitHub 拒绝了此操作:你的账号没有权限。";
if (t === "GITHUB_NOT_FOUND") return "找不到目标(可能已删除或仓库不可访问)。";
return `操作失败:${t}`;
}
async function cmdLogin(env: Env, msg: TelegramMessage): Promise<void> {
const chatId = chatIdOf(msg);
const telegramUserId = userIdOf(msg);
if (!chatId || !telegramUserId) return;
const topicId = msg.message_thread_id != null ? String(msg.message_thread_id) : undefined;
const clientId = env.GITHUB_CLIENT_ID;
if (!clientId) return reply(env, chatId, topicId, "服务器未配置 GitHub OAuthGITHUB_CLIENT_ID。");
const state = crypto.randomUUID().replace(/-/g, "");
await env.KV.put(
`state:${state}`,
JSON.stringify({
redirectTo: "/",
telegramUserId,
telegramChatId: chatId,
expiresAt: Date.now() + 600_000,
}),
{ expirationTtl: 600 },
);
const url = getOAuthURL(clientId, state);
await reply(env, chatId, topicId, `点击链接授权 GitHub即可用**本人身份**评论10 分钟内有效):\n${url}`);
}
async function cmdLogout(env: Env, msg: TelegramMessage): Promise<void> {
const chatId = chatIdOf(msg);
const telegramUserId = userIdOf(msg);
if (!chatId || !telegramUserId) return;
const topicId = msg.message_thread_id != null ? String(msg.message_thread_id) : undefined;
await removeTelegramLink(env.DB, telegramUserId);
await reply(env, chatId, topicId, "已解绑你的 GitHub 账号。");
}
async function cmdComment(env: Env, msg: TelegramMessage, body: string): Promise<void> {
const chatId = chatIdOf(msg);
const telegramUserId = userIdOf(msg);
if (!chatId || !telegramUserId) return;
const topicId = msg.message_thread_id != null ? String(msg.message_thread_id) : undefined;
const githubUserId = await getTelegramLink(env.DB, telegramUserId);
if (!githubUserId) {
return reply(env, chatId, topicId, "你还没有绑定 GitHub 账号,请先使用 /gh login。");
}
const source = msg.reply_to_message;
const target = source ? extractTarget(source) : null;
if (!target) {
return reply(
env,
chatId,
topicId,
"找不到 issue / PR 链接,请在对应的 GitHub 通知消息上回复 /gh comment。",
);
}
try {
const { htmlUrl, login } = await commentAsUser(
env.KV,
githubUserId,
target.owner,
target.repo,
target.number,
body,
);
await reply(env, chatId, topicId, `已以 **@${login}** 身份评论:${htmlUrl}`);
} catch (err) {
await reply(env, chatId, topicId, errText(err));
}
}
async function cmdMergeClose(env: Env, msg: TelegramMessage, op: "merge" | "close"): Promise<void> {
const chatId = chatIdOf(msg);
const telegramUserId = userIdOf(msg);
if (!chatId || !telegramUserId) return;
const topicId = msg.message_thread_id != null ? String(msg.message_thread_id) : undefined;
const githubUserId = await getTelegramLink(env.DB, telegramUserId);
if (!githubUserId) {
return reply(env, chatId, topicId, "你还没有绑定 GitHub 账号,请先使用 /gh login。");
}
const source = msg.reply_to_message;
const target = source ? extractTarget(source, true) : null;
if (!target) {
return reply(
env,
chatId,
topicId,
"找不到 PR 链接,请在对应的 GitHub PR 通知消息上回复 /gh merge 或 /gh close。",
);
}
try {
if (op === "merge") {
await mergePullRequestAsUser(env.KV, githubUserId, target.owner, target.repo, target.number);
} else {
await closePullRequestAsUser(env.KV, githubUserId, target.owner, target.repo, target.number);
}
const label = op === "merge" ? "合并" : "关闭";
await reply(env, chatId, topicId, `✅ 已${label} PR ${target.owner}/${target.repo}#${target.number}`);
} catch (err) {
await reply(env, chatId, topicId, errText(err));
}
}
/** Register the Telegram webhook to point at this worker (called from cron). */
export async function syncTelegramWebhook(env: Env): Promise<void> {
const token = env.TELEGRAM_TOKEN;
if (!token) return;
const baseUrl = env.BASE_URL;
if (!baseUrl) return;
const secret = env.TELEGRAM_WEBHOOK_SECRET;
const res = await fetch(`https://api.telegram.org/bot${token}/setWebhook`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({
url: `${baseUrl.replace(/\/$/, "")}/telegram/webhook`,
secret_token: secret || undefined,
allowed_updates: ["message"],
}),
});
if (!res.ok) {
const err = await res.text();
log.error({ status: res.status, err }, "Failed to set Telegram webhook");
}
}
/** Handle a single Telegram update (message). Called from the webhook route. */
export async function handleTelegramUpdate(env: Env, update: unknown): Promise<void> {
const message = (update as { message?: TelegramMessage })?.message;
if (!message?.text) return;
const text = message.text.trim();
const m = text.match(/^\/gh(?:\s+|$)(.*)$/s);
if (!m) return;
const rest = m[1].trim();
const [sub, ...args] = rest.split(/\s+/);
const body = args.join(" ").trim();
switch (sub) {
case "login":
return cmdLogin(env, message);
case "logout":
return cmdLogout(env, message);
case "comment":
if (!body) {
const chatId = chatIdOf(message);
const topicId =
message.message_thread_id != null ? String(message.message_thread_id) : undefined;
if (chatId) await reply(env, chatId, topicId, "请附上评论内容:/gh comment 你的评论");
return;
}
return cmdComment(env, message, body);
case "merge":
return cmdMergeClose(env, message, "merge");
case "close":
return cmdMergeClose(env, message, "close");
default:
log.info({ text }, "Unhandled /gh command from Telegram");
}
}

View file

@ -1,10 +1,17 @@
import type { Route, Env, NeutralMessage } from "../../types";
import type { RouteTarget, Env, NeutralMessage } from "../../types";
import type { PlatformDriver, SendResult } from "../types";
import { sendMessage } from "./rest";
import { renderNeutralMessage } from "./render";
export class TelegramDriver implements PlatformDriver {
readonly id = "telegram";
async send(_message: NeutralMessage, _target: Route["target"], _env: Env): Promise<SendResult> {
return { ok: false, error: "Telegram driver not implemented yet" };
async send(message: NeutralMessage, target: RouteTarget, env: Env): Promise<SendResult> {
const chatId = target.chatId ?? "";
if (!chatId) {
return { ok: false, error: "target.chatId is required", errorCode: "NO_TARGET" };
}
const token = env.TELEGRAM_TOKEN ?? "";
return sendMessage(token, chatId, renderNeutralMessage(message), target.topicId);
}
}

View file

@ -0,0 +1,45 @@
import type { NeutralMessage } from "../../types";
function esc(s: string): string {
return s
.replace(/&/g, "&amp;")
.replace(/</g, "&lt;")
.replace(/>/g, "&gt;")
.replace(/"/g, "&quot;");
}
function inlineUrl(url?: string, text?: string): string {
const label = esc(text ?? url ?? "");
if (!url) return label;
return `<a href="${esc(url)}">${label}</a>`;
}
export function renderNeutralMessage(message: NeutralMessage): string {
const parts: string[] = [];
const title = inlineUrl(message.url, message.title);
parts.push(`<b>${title}</b>`);
if (message.author) {
const author = message.author.url
? `<a href="${esc(message.author.url)}">${esc(message.author.name)}</a>`
: esc(message.author.name);
parts.push(`👤 ${author}`);
}
if (message.description) {
parts.push(esc(message.description));
}
for (const field of message.fields ?? []) {
parts.push(`<b>${esc(field.name)}</b>: ${esc(field.value)}`);
}
const meta: string[] = [];
if (message.footer) meta.push(esc(message.footer));
if (message.timestamp) meta.push(esc(message.timestamp));
if (meta.length > 0) {
parts.push(`<i>${meta.join(" · ")}</i>`);
}
return parts.join("\n");
}

View file

@ -0,0 +1,82 @@
import { log } from "../../lib/log";
import type { SendResult } from "../types";
const TELEGRAM_API = "https://api.telegram.org";
interface TelegramResponse {
ok?: boolean;
description?: string;
result?: { message_id?: number };
}
export async function sendMessage(
token: string,
chatId: string,
text: string,
topicId?: string,
): Promise<SendResult> {
if (!token) {
return { ok: false, error: "TELEGRAM_TOKEN not configured", errorCode: "NO_TOKEN" };
}
const body: Record<string, unknown> = {
chat_id: chatId,
text,
parse_mode: "HTML",
disable_web_page_preview: true,
};
if (topicId) {
body.message_thread_id = Number(topicId);
}
const url = `${TELEGRAM_API}/bot${token}/sendMessage`;
let lastStatus = 0;
let lastError = "";
for (let attempt = 0; attempt < 3; attempt++) {
try {
const res = await fetch(url, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify(body),
});
lastStatus = res.status;
const data = (await res.json().catch(() => null)) as TelegramResponse | null;
if (res.status === 429) {
const retryAfter = (data as { retry_after?: number })?.retry_after ?? 1;
lastError = data?.description ?? `Rate limited (retry_after=${retryAfter})`;
await new Promise((r) => setTimeout(r, retryAfter * 1000));
continue;
}
if (!res.ok) {
lastError = data?.description ?? `HTTP ${res.status}`;
log.error({ status: res.status, err: lastError, chatId, attempts: attempt + 1 }, "Telegram API error");
return {
ok: false,
error: lastError,
errorCode: res.status >= 500 ? "TELEGRAM_5XX" : "TELEGRAM_ERROR",
status: res.status,
attempts: attempt + 1,
};
}
return {
ok: true,
status: res.status,
messageId: data?.result?.message_id != null ? String(data.result.message_id) : undefined,
attempts: attempt + 1,
};
} catch (err) {
lastError = err instanceof Error ? err.message : String(err);
log.error({ err, chatId, attempts: attempt + 1 }, "Failed to send Telegram message");
if (attempt === 2) {
return { ok: false, error: lastError, errorCode: "NETWORK", status: lastStatus, attempts: attempt + 1 };
}
await new Promise((r) => setTimeout(r, 500 * (attempt + 1)));
}
}
return { ok: false, error: lastError || "Failed to send Telegram message", errorCode: "RETRIES", status: lastStatus, attempts: 3 };
}

View file

@ -0,0 +1,45 @@
import type { Env } from "../../types";
import { handleTelegramUpdate } from "./commands";
const MAX_BODY_SIZE = 1024 * 1024;
function timingSafeEqual(a: string, b: string): boolean {
if (a.length !== b.length) return false;
let diff = 0;
for (let i = 0; i < a.length; i++) {
diff |= a.charCodeAt(i) ^ b.charCodeAt(i);
}
return diff === 0;
}
/** Handle a POST to the Telegram webhook endpoint. */
export async function handleTelegramWebhookRequest(request: Request, env: Env): Promise<Response> {
const contentLength = Number(request.headers.get("content-length") ?? 0);
if (contentLength > MAX_BODY_SIZE) {
return new Response("Request too large", { status: 413 });
}
const secret = env.TELEGRAM_WEBHOOK_SECRET;
if (secret) {
const token = request.headers.get("X-Telegram-Bot-Api-Secret-Token");
if (!token || !timingSafeEqual(token, secret)) {
return new Response("Invalid secret", { status: 401 });
}
}
const rawBody = await request.text();
if (rawBody.length > MAX_BODY_SIZE) {
return new Response("Request too large", { status: 413 });
}
let update: unknown;
try {
update = JSON.parse(rawBody);
} catch {
return new Response("Invalid JSON", { status: 400 });
}
// Telegram expects a quick 200; process commands in the background.
await handleTelegramUpdate(env, update);
return new Response("ok");
}

View file

@ -1,4 +1,4 @@
import type { Route, Env, NeutralMessage } from "../types";
import type { RouteTarget, Env, NeutralMessage } from "../types";
export interface SendResult {
ok: boolean;
@ -11,5 +11,5 @@ export interface SendResult {
export interface PlatformDriver {
readonly id: string;
send(message: NeutralMessage, target: Route["target"], env: Env): Promise<SendResult>;
send(message: NeutralMessage, target: RouteTarget, env: Env): Promise<SendResult>;
}

View file

@ -97,3 +97,29 @@ export async function getDiscordLink(
export async function removeDiscordLink(db: D1Database, discordUserId: string): Promise<void> {
await db.prepare("DELETE FROM discord_links WHERE discord_user_id = ?").bind(discordUserId).run();
}
export async function saveTelegramLink(
db: D1Database,
telegramUserId: string,
githubUserId: string,
): Promise<void> {
await db
.prepare("INSERT OR REPLACE INTO telegram_links (telegram_user_id, github_user_id) VALUES (?, ?)")
.bind(telegramUserId, githubUserId)
.run();
}
export async function getTelegramLink(
db: D1Database,
telegramUserId: string,
): Promise<string | null> {
const { results } = await db
.prepare("SELECT github_user_id FROM telegram_links WHERE telegram_user_id = ?")
.bind(telegramUserId)
.all<{ github_user_id: string }>();
return results[0]?.github_user_id ?? null;
}
export async function removeTelegramLink(db: D1Database, telegramUserId: string): Promise<void> {
await db.prepare("DELETE FROM telegram_links WHERE telegram_user_id = ?").bind(telegramUserId).run();
}

View file

@ -1,5 +1,6 @@
import { createServer } from "./server";
import { syncCommands } from "./drivers/discord/commands";
import { syncTelegramWebhook } from "./drivers/telegram/commands";
import type { Env } from "./types";
import { log } from "./lib/log";
@ -16,5 +17,10 @@ export default {
} catch (err) {
log.error({ err }, "Discord command sync from cron failed");
}
try {
await syncTelegramWebhook(env);
} catch (err) {
log.error({ err }, "Telegram webhook sync from cron failed");
}
},
};

View file

@ -4,6 +4,7 @@ import { verifySignature } from "./events/verify";
import { parseEvent } from "./events/parse";
import { dispatchEvent } from "./core/dispatch";
import { handleInteractionRequest } from "./drivers/discord/interactions";
import { handleTelegramWebhookRequest } from "./drivers/telegram/updates";
import { createOAuthRoutes } from "./web/oauth-routes";
import { createActionRoutes } from "./web/action-routes";
import { createAdminRoutes } from "./web/admin-routes";
@ -71,6 +72,7 @@ export function createServer(): Hono<{ Bindings: Env }> {
});
app.post("/discord/interactions", (c) => handleInteractionRequest(c.req.raw, c.env));
app.post("/telegram/webhook", (c) => handleTelegramWebhookRequest(c.req.raw, c.env));
app.notFound((c) => {
if (c.env.ASSETS) {

View file

@ -13,6 +13,8 @@ export interface Env {
GITHUB_REPO_URL?: string;
DISCORD_PUBLIC_KEY?: string;
DISCORD_APPLICATION_ID?: string;
TELEGRAM_TOKEN?: string;
TELEGRAM_WEBHOOK_SECRET?: string;
ASSETS?: Fetcher;
KV: KVNamespace;
DB: D1Database;
@ -33,15 +35,20 @@ export interface Config {
routes: Route[];
}
export interface RouteTarget {
platform?: "discord" | "telegram";
channelId?: string;
threadId?: string;
chatId?: string;
topicId?: string;
}
export interface Route {
id: string;
name: string;
enabled: boolean;
filters: Filter[];
target: {
channelId: string;
threadId?: string;
};
targets: RouteTarget[];
lang?: string;
groupId?: string;
/**

View file

@ -98,16 +98,66 @@ function validateRoutes(
return { ok: false, error: `route "${r.id}" filter[${j}].exclude must be boolean` };
}
}
const target = r.target as Record<string, unknown> | undefined;
if (!target || typeof target !== "object")
return { ok: false, error: `route "${r.id}" needs a target` };
const rawTarget = r.target as Record<string, unknown> | undefined;
const rawTargets = r.targets as unknown;
if (rawTargets === undefined && rawTarget && typeof rawTarget === "object") {
const legacy = validateTarget(r, rawTarget);
if (!legacy.ok) return legacy;
(r as Record<string, unknown>).targets = [legacy.target];
delete (r as Record<string, unknown>).target;
} else if (Array.isArray(rawTargets)) {
if (rawTargets.length === 0) {
return { ok: false, error: `route "${r.id}" needs at least one target` };
}
const normalized: Route["targets"] = [];
for (let j = 0; j < rawTargets.length; j++) {
const t = rawTargets[j] as Record<string, unknown>;
if (!t || typeof t !== "object") {
return { ok: false, error: `route "${r.id}".targets[${j}] is not an object` };
}
const result = validateTarget(r, t);
if (!result.ok) return result;
normalized.push(result.target);
}
(r as Record<string, unknown>).targets = normalized;
} else {
return { ok: false, error: `route "${r.id}" needs a targets array` };
}
}
return { ok: true, routes: routes as Route[] };
}
function validateTarget(
r: Record<string, unknown>,
target: Record<string, unknown>,
): { ok: true; target: Route["targets"][number] } | { ok: false; error: string } {
const platform = target.platform === undefined ? "discord" : target.platform;
if (platform !== "discord" && platform !== "telegram") {
return { ok: false, error: `route "${r.id}".target.platform must be "discord" or "telegram"` };
}
if (platform === "telegram") {
if (typeof target.chatId !== "string" || target.chatId.trim().length === 0)
return { ok: false, error: `route "${r.id}".target.chatId is required` };
if (target.topicId !== undefined && typeof target.topicId !== "string") {
return { ok: false, error: `route "${r.id}".target.topicId must be a string` };
}
} else {
if (typeof target.channelId !== "string" || target.channelId.trim().length === 0)
return { ok: false, error: `route "${r.id}".target.channelId is required` };
if (target.threadId !== undefined && typeof target.threadId !== "string") {
return { ok: false, error: `route "${r.id}".target.threadId must be a string` };
}
}
return { ok: true, routes: routes as Route[] };
return {
ok: true,
target: {
platform,
channelId: platform === "telegram" ? undefined : (target.channelId as string),
threadId: platform === "telegram" ? undefined : ((target.threadId as string) ?? undefined),
chatId: platform === "telegram" ? (target.chatId as string) : undefined,
topicId: platform === "telegram" ? ((target.topicId as string) ?? undefined) : undefined,
},
};
}
function validateGroups(

View file

@ -1,14 +1,17 @@
import { Hono } from "hono";
import { getOAuthURL, handleOAuthCallback } from "../github/oauth";
import { removeToken, saveDiscordLink } from "../github/store";
import { removeToken, saveDiscordLink, saveTelegramLink } from "../github/store";
import { createAdminSession, adminCookie } from "./session";
import { loadGroups, resolveScope, hasAnyAccess } from "./groups";
import { sendMessage } from "../drivers/telegram/rest";
import type { Env } from "../types";
interface PendingState {
redirectTo: string;
expiresAt: number;
discordUserId?: string;
telegramUserId?: string;
telegramChatId?: string;
}
function linkedPage(login: string): string {
@ -91,6 +94,19 @@ export function createOAuthRoutes(): Hono<{ Bindings: Env }> {
return c.json({ ok: true, discordUserId: pending.discordUserId, login: result.login });
}
// Telegram account-linking flow: bind the Telegram user to this GitHub account.
if (pending.telegramUserId) {
await saveTelegramLink(c.env.DB, pending.telegramUserId, result.userId);
if (pending.telegramChatId && c.env.TELEGRAM_TOKEN) {
await sendMessage(
c.env.TELEGRAM_TOKEN,
pending.telegramChatId,
`✅ GitHub 账号已绑定:**@${result.login}**。现在可以用 /gh comment 评论了。`,
).catch(() => undefined);
}
return c.json({ ok: true, telegramUserId: pending.telegramUserId, login: result.login });
}
const isBrowser = (c.req.header("accept") ?? "").includes("text/html");
if (isBrowser) {
const groups = await loadGroups(c.env.KV);