Files
blog/blog-admin/src/lib/rss/cron.ts
T

407 lines
13 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// 抓取主流程:遍历订阅源 → 抓取 → 解析 → 去重 → 通知 → 写 KV
// 迁自 check-feeds.js 的 resolveFeed + main,作为 CF Cron 的 scheduled handler。
import type { Env } from '../../types';
import { fetchUrl } from './fetch';
import { discoverFeedUrl, tryCommonFeedPaths } from './discover';
import {
parseFeedXml,
parseJsonFeed,
parseCustomJson,
isJsonResponse,
type Article,
type FeedConfig,
type ParsedFeed,
} from './parse';
import { sendFeishuNotification } from './notify';
import { kvGetJson, kvPutJson } from './util';
const MAX_SEEN = 500;
/** 单源抓取超时(原来 30s:一个挂掉的源能把整批拖到 4 分钟以上) */
const FEED_TIMEOUT = 8000;
/** 连续失败这么多次就进入退避期,不再每次抓(否则每小时白等 8 秒) */
const FAIL_THRESHOLD = 3;
/** 退避时长:6 小时(期间只在其它源都抓完后才可能被跳过) */
const FAIL_BACKOFF_MS = 6 * 3600 * 1000;
interface FailRecord { n: number; at: number }
/** 抓取 + 解析单个订阅源(含 feed 自动发现)。失败返回 null。 */
async function resolveFeed(env: Env, feedConfig: FeedConfig): Promise<ParsedFeed | null> {
const url = feedConfig.url;
const format = feedConfig.format || 'auto';
const forceProxy = feedConfig.proxy === true;
let response;
try {
response = await fetchUrl(env, url, FEED_TIMEOUT, forceProxy);
} catch {
// 直连 + 代理都失败:不再重复抓同一个地址(原来会白等一个超时),
// 直接把常见的 feed 路径猜一遍就放弃。
response = undefined;
}
if (!response) {
const guessedUrl = await tryCommonFeedPaths(env, url);
if (!guessedUrl) return null;
try {
response = await fetchUrl(env, guessedUrl, FEED_TIMEOUT, forceProxy);
} catch {
return null;
}
}
const { text, contentType } = response;
const isJson = isJsonResponse(text, contentType);
if (format === 'json' || (format === 'auto' && isJson)) {
try {
const json = JSON.parse(text);
if (
json.version?.includes('jsonfeed.org') ||
(json.items && Array.isArray(json.items) && !feedConfig.path)
) {
return parseJsonFeed(json, url);
}
return parseCustomJson(json, feedConfig);
} catch {
return null;
}
}
if (
text.includes('<rss') ||
text.includes('<feed') ||
text.includes('<channel') ||
text.includes('<entry')
) {
return parseFeedXml(text, url);
}
// HTML 发现
const feedUrl = discoverFeedUrl(text, url);
if (feedUrl) {
const feedResp = await fetchUrl(env, feedUrl, FEED_TIMEOUT, forceProxy);
if (feedResp) {
if (isJsonResponse(feedResp.text, feedResp.contentType)) {
try {
const json = JSON.parse(feedResp.text);
if (json.version?.includes('jsonfeed.org') || json.items) return parseJsonFeed(json, feedUrl);
return parseCustomJson(json, { url: feedUrl });
} catch {
return null;
}
}
return parseFeedXml(feedResp.text, feedUrl);
}
}
const guessedUrl = await tryCommonFeedPaths(env, url);
if (guessedUrl) {
const guessResp = await fetchUrl(env, guessedUrl, FEED_TIMEOUT, forceProxy);
if (guessResp) {
if (isJsonResponse(guessResp.text, guessResp.contentType)) {
try {
const json = JSON.parse(guessResp.text);
if (json.version?.includes('jsonfeed.org') || json.items) return parseJsonFeed(json, guessedUrl);
return parseCustomJson(json, { url: guessedUrl });
} catch {
return null;
}
}
return parseFeedXml(guessResp.text, guessedUrl);
}
}
return null;
}
/** 从远程 JSON / KV 加载订阅源列表 */
async function loadFeeds(env: Env): Promise<FeedConfig[]> {
// 优先从 KV 读(管理页维护的最新订阅源)
const kvFeeds = await kvGetJson<{ feeds: FeedConfig[] }>(env, 'feeds_config', { feeds: [] });
if (kvFeeds.feeds && kvFeeds.feeds.length > 0) {
return kvFeeds.feeds.map((f) => (typeof f === 'string' ? { url: f } : f));
}
// 回退:FEEDS_URL 远程 JSON
if (env.FEEDS_URL) {
try {
const res = await fetch(env.FEEDS_URL, {
headers: { 'User-Agent': 'Mozilla/5.0 (compatible; RSSBot/1.0)' },
});
if (res.ok) {
const json = (await res.json()) as any;
if (Array.isArray(json)) {
return json.map((item: string | FeedConfig) =>
typeof item === 'string' ? { url: item } : item,
);
}
if (json.feeds && Array.isArray(json.feeds)) {
return json.feeds.map((item: string | FeedConfig) =>
typeof item === 'string' ? { url: item } : item,
);
}
}
} catch {
// 忽略,走空
}
}
return [];
}
interface SeenData {
lastCheck: string | null;
articles: Record<string, string>;
}
/** scheduled 入口:每天抓取一次 */
/**
* 抓取一批订阅源。
* 免费版 Worker 每次调用限 50 个 subrequest,66 个源全量一把抓会炸,
* 所以按 offset/limit 分批轮转(一天两批跑完全部);latest 按 feed 维度合并写入。
*/
export interface CronStats {
/** 本批计划抓取的源数 */
batch: number;
/** 订阅源总数 */
total: number;
ok: number;
failed: number;
/** 本批抓到的文章条数(去重前) */
articles: number;
/** 判定为新文章的条数 */
newArticles: number;
/** 是否 dry-run(不通知、不写 KV) */
dry: boolean;
durationMs: number;
/** 因连续失败处于退避期、本次跳过的源数 */
skipped: number;
/** 最慢的几个源,用于排查拖后腿的 feed */
slowest: { url: string; ms: number }[];
/** 本次失败的源(最多列 10 个,便于排查) */
failedUrls: string[];
}
export async function runCron(
env: Env,
opts?: { offset?: number; limit?: number; dry?: boolean; rotate?: boolean },
): Promise<CronStats> {
const startedAt = Date.now();
const dry = opts?.dry === true;
const rotate = opts?.rotate === true;
const offset = Math.max(opts?.offset ?? 0, 0);
const batchLimit = Math.max(opts?.limit ?? 0, 0);
const timings: { url: string; ms: number }[] = [];
let okCount = 0;
let failedCount = 0;
let articleCount = 0;
let skippedCount = 0;
const failedUrls: string[] = [];
const emptyStats: CronStats = {
batch: 0, total: 0, ok: 0, failed: 0, articles: 0, newArticles: 0,
dry, durationMs: 0, skipped: 0, slowest: [], failedUrls: [],
};
console.log('🔍 开始检查 RSS Feed...');
const allFeeds = await loadFeeds(env);
if (allFeeds.length === 0) {
console.log('⚠️ 没有可用的订阅源');
return emptyStats;
}
// 失败退避表:连续失败 >= 阈值的源,在退避期内不再每次白等超时
const failMap = await kvGetJson<Record<string, FailRecord>>(env, 'rss_fail', {});
const nowTs = Date.now();
const inBackoff = (u: string): boolean => {
const r = failMap[u];
return !!r && r.n >= FAIL_THRESHOLD && nowTs - r.at < FAIL_BACKOFF_MS;
};
const markFail = (u: string): void => {
const r = failMap[u];
failMap[u] = { n: (r?.n ?? 0) + 1, at: nowTs };
};
const clearFail = (u: string): void => {
if (failMap[u]) delete failMap[u];
};
console.log(`📋 共 ${allFeeds.length} 个订阅源`);
for (const u of Object.keys(failMap)) if (!inBackoff(u)) delete failMap[u]; // 退避期满自动重试
// 环形取本批:offset 开始取 limit 个(limit=0 表示全量)。
// rotate=true 时从 KV 里的游标继续(定时任务用,保证每轮都能覆盖到所有源)。
const feeds: FeedConfig[] = [];
let cursorAfter = offset;
if (batchLimit > 0 && allFeeds.length > batchLimit) {
let start = offset % allFeeds.length;
if (rotate && !dry) {
const saved = await env.RSS_KV.get('rss_cursor');
const n = parseInt(saved || '0', 10);
if (Number.isFinite(n)) start = ((n % allFeeds.length) + allFeeds.length) % allFeeds.length;
}
let examined = 0;
let i = start;
while (feeds.length < batchLimit && examined < allFeeds.length) {
const f = allFeeds[i % allFeeds.length];
i += 1;
examined += 1;
if (inBackoff(f.url)) {
skippedCount += 1;
continue;
}
feeds.push(f);
}
cursorAfter = i % allFeeds.length;
if (rotate && !dry) await env.RSS_KV.put('rss_cursor', String(cursorAfter));
console.log(
`📋 本批 ${feeds.length} 个(cursor=${start}${skippedCount ? `,跳过退避中 ${skippedCount} 个` : ''})`,
);
} else {
feeds.push(...allFeeds);
}
const seenData = await kvGetJson<SeenData>(env, 'seen_articles', {
lastCheck: null,
articles: {},
});
const allNewArticles: Article[] = [];
const allFetchedArticles: (Article & { siteUrl: string })[] = [];
for (const feedConfig of feeds) {
const url = feedConfig.url;
console.log(`🔍 检查: ${url}`);
const t0 = Date.now();
try {
const result = await resolveFeed(env, feedConfig);
timings.push({ url, ms: Date.now() - t0 });
if (!result) {
failedCount++;
markFail(url);
if (failedUrls.length < 10) failedUrls.push(url);
console.log(' ❌ 未找到 Feed');
continue;
}
okCount++;
articleCount += result.articles.length;
clearFail(url);
const { feedTitle, articles } = result;
console.log(` 📰 ${feedTitle || '未知'} - 共 ${articles.length} 篇`);
const recentArticles = articles.slice(0, 10);
allFetchedArticles.push(
...recentArticles.map((a) => ({ ...a, feedTitle: feedTitle || '未知', siteUrl: url })),
);
const newOnes = recentArticles.filter((a) => !seenData.articles[a.link]);
if (newOnes.length > 0) {
console.log(` 🆕 发现 ${newOnes.length} 篇新文章`);
allNewArticles.push(...newOnes);
} else {
console.log(' ✅ 无新文章');
}
for (const a of recentArticles) {
seenData.articles[a.link] = a.pubDate;
}
} catch (err) {
timings.push({ url, ms: Date.now() - t0 });
failedCount++;
markFail(url);
if (failedUrls.length < 10) failedUrls.push(url);
console.error(` ❌ 抓取失败: ${(err as Error).message}`);
}
}
// 通知(仅最近 24h 内的)
const oneDayAgo = Date.now() - 86400000;
const recentNew = allNewArticles.filter(
(a) => new Date(a.pubDate).getTime() > oneDayAgo || !seenData.lastCheck,
);
if (dry) {
console.log(`🧪 dry-run:跳过通知(本应通知 ${recentNew.length} 篇)与 KV 写入`);
} else if (recentNew.length > 0) {
await sendFeishuNotification(env, recentNew);
} else if (allNewArticles.length > 0) {
console.log('ℹ️ 发现新文章但超过 24 小时,不推送');
} else {
console.log('✅ 所有博客均无新文章');
}
// 写 latest 缓存(供 /api/results、/api/articles 读)
// 分批模式:本批源的条目替换旧缓存里的同源条目,其他源的保留
const byFeed: Record<string, { articles: Article[]; siteUrl: string }> = {};
for (const a of allFetchedArticles) {
const key = a.feedTitle || '未知';
if (!byFeed[key]) byFeed[key] = { articles: [], siteUrl: a.siteUrl || '' };
byFeed[key].articles.push({ title: a.title, link: a.link, pubDate: a.pubDate, author: a.author, feedTitle: key });
}
if (!dry) {
try {
await kvPutJson(env, 'rss_fail', failMap);
} catch {
/* 失败表写失败不影响主流程 */
}
}
const durationMs = Date.now() - startedAt;
const stats: CronStats = {
batch: feeds.length,
total: allFeeds.length,
ok: okCount,
failed: failedCount,
articles: articleCount,
newArticles: allNewArticles.length,
dry,
durationMs,
skipped: skippedCount,
slowest: timings.sort((a, b) => b.ms - a.ms).slice(0, 5),
failedUrls,
};
if (dry) return stats;
const prev = await kvGetJson<{ timestamp: string | null; total: number; feeds: { name: string; siteUrl: string; favicon: string; articles: Article[] }[] }>(
env,
'latest',
{ timestamp: null, total: 0, feeds: [] },
);
const batchUrls = new Set(feeds.map((f) => f.url));
const kept = prev.feeds.filter((f) => !batchUrls.has(f.siteUrl || f.name));
const mergedFeeds = [
...kept,
...Object.entries(byFeed).map(([name, { articles, siteUrl }]) => ({
name,
siteUrl,
// 必须写绝对地址:友圈页在 usj.cc 上,相对路径 /api/* 会 404 → 退化成第三方默认头像
favicon: siteUrl
? `${env.PUBLIC_API_BASE || 'https://api.200181.xyz'}/api/favicon?url=${encodeURIComponent(siteUrl)}`
: '',
articles,
})),
];
const mergedTotal = mergedFeeds.reduce((n, f) => n + f.articles.length, 0);
await kvPutJson(env, 'latest', {
timestamp: new Date().toISOString(),
total: mergedTotal,
feeds: mergedFeeds,
});
// 清理 + 保存去重记录
const entries = Object.entries(seenData.articles);
if (entries.length > MAX_SEEN) {
entries.sort((a, b) => new Date(b[1]).getTime() - new Date(a[1]).getTime());
seenData.articles = Object.fromEntries(entries.slice(0, MAX_SEEN));
}
seenData.lastCheck = new Date().toISOString();
await kvPutJson(env, 'seen_articles', seenData);
console.log(`📝 已保存去重记录 (${Object.keys(seenData.articles).length} 条)`);
return stats;
}