iboard/scripts/fetchers/news.fetcher.js
cheney 6bacda05f8
All checks were successful
Docker Build and Push / build-image (push) Successful in 3m2s
fix bugs
2026-08-06 17:43:21 +08:00

137 lines
5.6 KiB
JavaScript

// 新闻 fetcher: 每天 3 点调用 mx_finance_search_news 拉国内外经济政治新闻
// 按影响大小排序, 取前 10 条 (LLM 可能多/少几条, 全部入 DB, 按 impactRank 标记)
// 数据直接写 News 表, 与 Indicator 解耦
import axios from 'axios'
import mysql from 'mysql2/promise'
const MCP_BASE_DEFAULT = 'https://mxapi.eastmoney.com/mxds/mcp'
const TIMEOUT = 90000
const TOP_N = 10
const QUERIES = [
'今天国内外最新重大经济政治新闻, 按影响大小排序, 给出前 ' + TOP_N + ' 条, 每条含标题、摘要、发布时间、来源、链接',
]
class McpClient {
constructor(baseUrl, apiKey) { this.baseUrl = baseUrl; this.apiKey = apiKey; this.id = 0 }
async _request(method, params) {
const payload = { jsonrpc: '2.0', id: ++this.id, method, ...(params ? { params } : {}) }
const resp = await axios.post(this.baseUrl, payload, {
timeout: TIMEOUT,
headers: { 'Content-Type': 'application/json', 'Accept': 'application/json, text/event-stream', 'em_api_key': this.apiKey },
responseType: 'text', transformResponse: [(d) => d], validateStatus: () => true,
})
if (resp.status >= 400) throw new Error('HTTP ' + resp.status + ': ' + String(resp.data).slice(0, 300))
const parsed = JSON.parse(String(resp.data || ''))
if (parsed.error) throw new Error('MCP error ' + parsed.error.code + ': ' + parsed.error.message)
return parsed.result || parsed
}
async callTool(name, args = {}) {
const r = await this._request('tools/call', { name, arguments: args })
if (r.isError) throw new Error('tool error: ' + (r.content || []).map((c) => c.text || '').join('; '))
return r
}
}
function parsePublishedAt(s) {
if (!s) return null
const d = new Date(String(s).replace(/-/g, '/'))
return isNaN(d.getTime()) ? null : d
}
function pickNewsSheet(sheets) {
if (!Array.isArray(sheets) || sheets.length === 0) return null
return sheets.find((s) => s.sheetName && (s.sheetName.includes('资讯') || s.sheetName.includes('新闻'))) || sheets[0]
}
// items[i] = [title, summary, publishedAt, source, url]
function parseNewsSheet(sheet, impactBase = 0) {
if (!sheet || !sheet.items) return []
const out = []
for (let i = 0; i < sheet.items.length; i++) {
const r = sheet.items[i]
const title = String(r[0] || '').trim().slice(0, 190)
if (!title) continue
const publishedAt = parsePublishedAt(r[2])
if (!publishedAt) continue
const body = String(r[1] || '').trim()
// 数据源段落间用全角空格连接, 入库前统一转成换行, 保证详情页有段落感
.replace(/\u3000+/g, '\n')
.replace(/\u00a0/g, ' ')
.replace(/[ \t]+\n/g, '\n')
.replace(/\n{3,}/g, '\n\n')
.trim()
out.push({
title,
summary: body.slice(0, 180),
content: body.slice(0, 8000), // 正文全文 (MCP "摘要"列实际是正文), 上限防异常超长
source: String(r[3] || '').trim().slice(0, 80),
url: String(r[4] || '').trim().slice(0, 190), // 列宽 varchar(191), 超长会被 MySQL 拒绝
publishedAt,
impactRank: impactBase + i + 1,
})
}
return out
}
export const META = { name: 'news-mcp', description: '妙想 MCP 新闻资讯 (mx_finance_search_news)' }
export async function fetchNewsList() {
const key = process.env.EM_API_KEY
if (!key) throw new Error('缺少 EM_API_KEY')
const base = process.env.EASTMONEY_MCP_BASE || MCP_BASE_DEFAULT
const client = new McpClient(base, key)
const all = []
let impactBase = 0
for (const q of QUERIES) {
try {
console.log('[News] 调用: ' + q)
const r = await client.callTool('mx_finance_search_news', { query: q })
const text = (r.content || []).map((c) => c.text || '').join('')
const payload = JSON.parse(text)
const sheets = payload.data || (Array.isArray(payload) ? payload : [payload])
const sheet = pickNewsSheet(sheets)
const news = parseNewsSheet(sheet, impactBase)
console.log('[News] 解析 ' + news.length + ' 条')
all.push(...news)
impactBase = all.length
} catch (e) {
console.error('[News] 拉取失败:', e.message)
}
}
return all
}
function parseUrl(url) {
const u = new URL(url)
return { host: u.hostname, port: Number(u.port) || 3306, user: decodeURIComponent(u.username), password: decodeURIComponent(u.password), database: u.pathname.replace(/^\//, '') }
}
export async function persistNews(newsList) {
if (!process.env.DATABASE_URL) throw new Error('缺少 DATABASE_URL')
if (!newsList || newsList.length === 0) return { inserted: 0, skipped: 0 }
const conn = await mysql.createConnection(parseUrl(process.env.DATABASE_URL))
let inserted = 0, skipped = 0, failed = 0
try {
for (const n of newsList) {
try {
// 注意: content 列依赖迁移 add_news_content 已应用, 发布顺序须先迁移后跑本脚本
// 重复命中 (title, publishedAt) 时只回填仍为 NULL 的 content, 避免覆盖已有全文
const [r] = await conn.execute(
'INSERT INTO `News` (`title`, `summary`, `content`, `source`, `url`, `publishedAt`, `impactRank`, `createdAt`) ' +
'VALUES (?, ?, ?, ?, ?, ?, ?, NOW(3)) ' +
'AS new ON DUPLICATE KEY UPDATE `impactRank` = new.`impactRank`, `content` = IFNULL(News.`content`, new.`content`)',
[n.title, n.summary, n.content || null, n.source, n.url, n.publishedAt, n.impactRank],
)
if (r.affectedRows === 1) inserted++
else skipped++
} catch (e) {
failed++
console.error('[News] 单条入库失败:', n.title, e.message)
}
}
} finally {
await conn.end()
}
return { inserted, skipped, failed }
}