// 新闻 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() 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 } }