定时任务
+ 任务管理
+
+ {{ drawerData.description }}
diff --git a/AGENTS.md b/AGENTS.md index 0d47861..cb2ffb0 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -4,7 +4,8 @@ - **Backend**: FastAPI 0.104 + SQLAlchemy 2.0 + PostgreSQL 15 (`yzr_nr`) - **Frontend**: Vue 3 (CDN, no build step) + Element Plus — static HTML served by FastAPI - **Auth**: JWT (`python-jose` + bcrypt), default admin `admin/admin123` -- **Scheduler**: APScheduler (daily cron: 01:30 collect, 02:30 sync, 03:30 generate, 04:30 optimize, 05:00 optimize_sources) +- **Scheduler**: APScheduler (daily cron: 01:00 searchcache, 01:10 trends, 01:30 collect, 02:00 generate, 03:00 optimize, 05:00 sources, 06:00 metrics) +- **Task DB**: `TaskLog` (module_id/status/error_trace/result_data/triggered_by) + `TaskConfig` (params/enabled/schedule) - **LLM**: Multi-provider (opencode-go primary, nvidia backup). API keys only in `.env`, not DB. ## Commands diff --git a/PROGRESS.md b/PROGRESS.md index c684eac..df89367 100644 --- a/PROGRESS.md +++ b/PROGRESS.md @@ -3,7 +3,7 @@ > 本文件为项目进度唯一真理源,所有进度信息以此为准。 > 其他文档中的进度描述一律以本文为准。 -**最后更新**:2026-05-18 (v16) +**最后更新**:2026-05-22 (v17) --- @@ -47,7 +47,7 @@ | 模块 | 状态 | 说明 | |------|------|------| | 后端 API (auth/topics/articles/publishing/calendar/metrics/assets/tasks/platform_config/admin) | ✅ 完成 | 核心 11 个 API 模块,JWT 认证 | -| 扩展 API (cases/audit/llm_configs/system_configs/optimizer_logs/task_logs) | ✅ 完成 | 新增案例库、审计、LLM 配置等模块 | +| 扩展 API (cases/audit/llm_configs/system_configs/optimizer_logs/task_logs/task_configs) | ✅ 完成 | 新增案例库、审计、LLM配置、任务配置、任务运行记录模块 | | 前端页面 (仪表盘/选题/日历/数据/素材/任务/平台/系统管理/用户/日志/文章) | ✅ 完成 | Vue 3 + Element Plus SPA | | 数据库 (PostgreSQL 15) | ✅ 运行中 | `yzr_nr` 库 | | 服务 | ✅ 运行中 | 端口 8001 | @@ -137,14 +137,22 @@ | metrics 冗余统计移除 | 2026-05-18 | 删除与仪表盘重复的 4 个概览卡片 | | writer max_tokens 提升 | 2026-05-19 | 标题/标签 max_tokens 500→1000,修复推理模型思考链占用导致输出截断 | | 预览滚动条修复 | 2026-05-19 | 改用 height:68vh 替代 flex+calc,避免 Element Plus 对话框内滚动冲突 | +| 定时任务 DB 化 | 2026-05-22 | TaskLog 新增 module_id/error_trace/triggered_by/result_data;TaskConfig 模型新建含 params/schedule/enabled;scheduler.py 全 7 个任务执行前后写 TaskRun;modules/status 从 DB 读取 | +| 平台配置弹窗修复 | 2026-05-22 | el-dialog 移入 tab 内部(与表格同级)解决响应式问题;表单字段补全 | +| AI 思考内容清洗 | 2026-05-22 | 新建 content_cleaner.py 集中管理清洗规则;writer.py/compliance_optimizer.py 统一引用;THINKING_PATTERNS 增强;strip_ai_preface 处理代码围栏块 | +| 复制正文去噪音 | 2026-05-22 | copyContent 只提取 p/h1-h4/li 元素,去 style/svg/script/img | +| platforms.html 入口合并 | 2026-05-22 | uni-nav 移除"平台"独立入口;admin.html 恢复"平台配置"tab 加启用中/全部筛选 | +| opencode_search.py 日志 | 2026-05-22 | 补 FileHandler + StreamHandler,解决管理后台显示"从未运行" | +| scheduler.json 导入修复 | 2026-05-22 | 补 import json,修复 sources(05:00) 执行时报错阻断 metrics(06:00) | ### ⏳ 待办 | 任务 | 优先级 | 备注 | |------|--------|------| +| admin.html 任务管理 tab(TaskConfig 参数编辑+TaskLog 历史时间轴) | 高 | 刚完成后端 DB 化,需完善前端 UI | | M4 第 4 篇文章发布 (首月目标) | 中 | 可用银发科技或 F01 补齐 | +| 现有文章重新创作(清除 AI 思考内容) | 中 | writer.py 已修复,新文章不会再有;旧文章需重跑 creator.py | | 数据追踪接入 (阅读量/互动) | 中 | 需要对接平台 API | -| 选题库数据库同步 | 低 | 将选题从 markdown 同步至 PG | | 归档 IMPLEMENTATION_PLAN.md | 低 | 内容已过时,与实际架构不符 | --- diff --git a/automation/data/sustainability_cases.json b/automation/data/sustainability_cases.json index bfadd0d..bd0f238 100644 --- a/automation/data/sustainability_cases.json +++ b/automation/data/sustainability_cases.json @@ -1,4 +1,140 @@ [ + { + "id": "SUS-6FC92B82", + "country": "Global", + "category": "循环消费", + "title": "氪星晚报|马斯克旗下SpaceX启动史上最大规模IPO计划;全国首张综合性司机服务地图上线;海关总署在粤发布《海关支持粤港澳大湾区建设若干措施》", + "core_idea": "
据美国联邦航空管理局(FA", + "data_facts": "数据点: 50, 40, 50", + "global_advantage": "需进一步分析全球优势", + "china_pain_point": "以旧换新流程繁琐、二手商品信任缺失、租赁市场不规范", + "localization_suggestion": "需基于中国现实调整实施", + "mvp_action": "建议先小规模试点验证", + "source_url": "https://36kr.com/p/3816942732133510?f=rss", + "credibility_rating": "⭐", + "china_applicability": "⭐⭐", + "collection_date": "2026-05-22", + "status": "待验证" + }, + { + "id": "SUS-5E0A0291", + "country": "Global", + "category": "循环消费", + "title": "新石器NewClaw:AI一体化解决方案,零门槛当无人车指挥官| 2026AI Partner·北京亦庄AI+产业大会", + "core_idea": "
\n\n管理一千台无人车需要多少人?答案是:一个人,一部手机,一句话。当自动驾驶逐渐“平权”,真正的瓶颈从技术转向了规模化运营。
\n
新石器用七年时间走完从合规落地、规模量产到万台运营的三级跳,如今推出AI Agent“Neo Claw”——让用户像聊天一样指挥车队,把单人管理效率从10台拉升到100台以上。颉晶华强调,A", + "data_facts": "数据点: 50, 1, 12", + "global_advantage": "需进一步分析全球优势", + "china_pain_point": "以旧换新流程繁琐、二手商品信任缺失、租赁市场不规范", + "localization_suggestion": "需基于中国现实调整实施", + "mvp_action": "建议先小规模试点验证", + "source_url": "https://36kr.com/p/3818927367046018?f=rss", + "credibility_rating": "⭐", + "china_applicability": "⭐⭐", + "collection_date": "2026-05-22", + "status": "待验证" + }, + { + "id": "SUS-47A4DB4F", + "country": "Global", + "category": "循环消费", + "title": "业绩快报 | 唯品会一季度净营收266亿元,SVIP用户贡献超50%线上销售额", + "core_idea": "
5月21日美股盘前,唯品会发布2025年第一季度财报。一季度内,唯品会实现净营收266亿元(人民币,下同),Non-GAAP净利润23亿元。
\n观察核心运营数据,一季度其实现商品交易总额(GMV)569亿元,同比增长8.6%,订单量1.73亿单,同比增长3.2%。同时,平台于该季度活跃用户数为4170万,同比实现正增长。
\n在春节期间穿戴和年货需求集中释放、透", + "data_facts": "数据点: 8.6, 3.2, 15", + "global_advantage": "需进一步分析全球优势", + "china_pain_point": "以旧换新流程繁琐、二手商品信任缺失、租赁市场不规范", + "localization_suggestion": "需基于中国现实调整实施", + "mvp_action": "建议先小规模试点验证", + "source_url": "https://36kr.com/p/3818915823764610?f=rss", + "credibility_rating": "⭐", + "china_applicability": "⭐⭐", + "collection_date": "2026-05-22", + "status": "待验证" + }, + { + "id": "SUS-8064DDFA", + "country": "Global", + "category": "循环消费", + "title": "从概念到产线一:AI在工业制造领域的深水区探索| 2026AI Partner·北京亦庄AI+产业大会", + "core_idea": "
\n\nAI在工业制造领域,不是“锦上添花”的辅助工具,而是“重新设计工厂”的核心引擎。从AOI报废板的秒级识别,到刀具参数的动态优化,再到打通设计、生产、供应链的全链路智能——这场对话告诉我们,每1%的效率提升都是真金白银,AI的价值不是叠加功能,而是把“人等货”变成“货等人”。
\n
演讲拆解了AI从概念到产线的落地路径", + "data_facts": "数据点: 1, 1, 80", + "global_advantage": "需进一步分析全球优势", + "china_pain_point": "以旧换新流程繁琐、二手商品信任缺失、租赁市场不规范", + "localization_suggestion": "需基于中国现实调整实施", + "mvp_action": "建议先小规模试点验证", + "source_url": "https://36kr.com/p/3818903062004870?f=rss", + "credibility_rating": "⭐", + "china_applicability": "⭐⭐", + "collection_date": "2026-05-22", + "status": "待验证" + }, + { + "id": "SUS-43DC92C1", + "country": "US", + "category": "循环消费", + "title": "城市级AI服务:从试点到常态化,机器人的实景作战与规模化落地| 2026AI Partner·北京亦庄AI+产业大会", + "core_idea": "
\n\n当Robotaxi还在为L4苦苦挣扎时,酷哇的环卫机器人、无人小巴、机器狗已经在50多个城市“上岗”赚钱了。
\n
具身智能最大的瓶颈不是算法,而是数据——没有量产就没有数据,没有数据就无法进化。酷哇的解法是“以战养战”:让机器人在真实运营中一边干活一边成长,用万台规模反哺模型迭代。李柯宏强调,中国是全球少有的支持机", + "data_facts": "数据点: 20, 20, 5500", + "global_advantage": "需进一步分析全球优势", + "china_pain_point": "以旧换新流程繁琐、二手商品信任缺失、租赁市场不规范", + "localization_suggestion": "需基于中国现实调整实施", + "mvp_action": "建议先小规模试点验证", + "source_url": "https://36kr.com/p/3818889074557829?f=rss", + "credibility_rating": "⭐", + "china_applicability": "⭐⭐", + "collection_date": "2026-05-22", + "status": "待验证" + }, + { + "id": "SUS-82A67B8B", + "country": "Global", + "category": "循环消费", + "title": "把确定性,写进农业:四个外行、两次失败、三千万学费换来的答案| 2026AI Partner·北京亦庄AI+产业大会", + "core_idea": "
\n\n两次失败、三千万学费——创业没有爽剧剧本,这是陆渔科技深耕农业AI真实的“入场券”。当99%的人还在用AI写文案、做设计时,有人把它扔进了鱼塘,只为解决一个最朴素的问题:不确定性。
\n
鲁敏用18年IT男转型“新农民”的经历,揭开了水产养殖最残酷的真相:1.38万亿的市场,数字化渗透率不足5%,一叶方塘,百万归零。", + "data_facts": "数据点: 99, 5, 300", + "global_advantage": "需进一步分析全球优势", + "china_pain_point": "以旧换新流程繁琐、二手商品信任缺失、租赁市场不规范", + "localization_suggestion": "需基于中国现实调整实施", + "mvp_action": "建议先小规模试点验证", + "source_url": "https://36kr.com/p/3818800487679111?f=rss", + "credibility_rating": "⭐", + "china_applicability": "⭐⭐", + "collection_date": "2026-05-22", + "status": "待验证" + }, + { + "id": "SUS-20916619", + "country": "Global", + "category": "循环消费", + "title": "从算力到价值:AI时代的基础设施重构与产业增长新引擎| 2026AI Partner·北京亦庄AI+产业大会", + "core_idea": "
\n\ntoken经济如何重塑AI产业链?从芯片到智算中心,从模型服务到终端应用,token正在成为贯穿全链的计价单位。而当推理算力需求超越训练,智算中心的角色正从算力仓库变为token工厂,一个万亿级市场的大门刚刚打开。
\n
token正在成为AI时代的新质生产力单位。宋琛指出,随着Agent成为新交互入口,单次任务to", + "data_facts": "数据点: 61, 60, 35", + "global_advantage": "需进一步分析全球优势", + "china_pain_point": "以旧换新流程繁琐、二手商品信任缺失、租赁市场不规范", + "localization_suggestion": "需基于中国现实调整实施", + "mvp_action": "建议先小规模试点验证", + "source_url": "https://36kr.com/p/3817502775329667?f=rss", + "credibility_rating": "⭐", + "china_applicability": "⭐⭐", + "collection_date": "2026-05-22", + "status": "待验证" + }, + { + "id": "SUS-37567F2B", + "country": "Global", + "category": "循环消费", + "title": "开场致辞 建设“全域人工智能之城” | 2026AI Partner·北京亦庄AI+产业大会", + "core_idea": "
\n\nAI的聚光灯正从炫酷的C端应用,转向轰鸣的工厂、无声的手术室和奔跑的人形机器人。当“落地”成为当下的关键词,2026北京亦庄AI+产业大会吹响了“走!带着AI去前线”的号角,在产业第一线求解人工智能的真实生产力。我们记录下这场务实者的聚会,捕捉那些让技术扎根泥土的坚定声音。
\n
5月19日,2026北京亦庄AI+产", + "data_facts": "数据点: 1.37, 40, 75", + "global_advantage": "需进一步分析全球优势", + "china_pain_point": "以旧换新流程繁琐、二手商品信任缺失、租赁市场不规范", + "localization_suggestion": "需基于中国现实调整实施", + "mvp_action": "建议先小规模试点验证", + "source_url": "https://36kr.com/p/3818784445465731?f=rss", + "credibility_rating": "⭐", + "china_applicability": "⭐⭐", + "collection_date": "2026-05-22", + "status": "待验证" + }, { "id": "SUS-0A6B894E", "country": "US", diff --git a/platform/backend/app/api/config_items.py b/platform/backend/app/api/config_items.py new file mode 100644 index 0000000..04b2236 --- /dev/null +++ b/platform/backend/app/api/config_items.py @@ -0,0 +1,390 @@ +from fastapi import APIRouter, Depends, HTTPException +from sqlalchemy.orm import Session +from typing import List, Optional +from pydantic import BaseModel, ConfigDict +from datetime import datetime +import json + +from ..database import get_db +from ..models import KeywordDomainMap, SensitiveWord, ContentCleanRule, TrendFieldMapping, CollectorCategory +from .auth import get_current_admin + +router = APIRouter(prefix="/api/admin/config", tags=["admin"]) + +DEFAULT_KEYWORD_DOMAIN_MAP = [ + {"pattern": r"AI|人工智能|大模型|GPT|机器学习|深度学习|聊天机器人|LLM", "domain": "AI工具", "sort_order": 1}, + {"pattern": r"远程|居家办公|自由职业|数字游民|远程协作", "domain": "远程工作", "sort_order": 2}, + {"pattern": r"可持续|环保|低碳|绿色|碳中和|循环|零浪费|垃圾分类|节能", "domain": "可持续生活", "sort_order": 3}, + {"pattern": r"知识管理|笔记|Obsidian|Notion|第二大脑|读书|阅读", "domain": "知识管理", "sort_order": 4}, + {"pattern": r"数字生活|数码|手机|电脑|智能|APP|应用、软件", "domain": "数字生活", "sort_order": 5}, + {"pattern": r"科技|人文|教育|心理|哲学|社会学", "domain": "科技人文", "sort_order": 6}, +] + +DEFAULT_TREND_FIELD_MAP = [ + {"trend_keyword": "远程工作", "field_name": "未来工作方式", "sort_order": 1}, + {"trend_keyword": "AI工具", "field_name": "AI与效率", "sort_order": 2}, + {"trend_keyword": "可持续生活", "field_name": "可持续生活系统", "sort_order": 3}, + {"trend_keyword": "知识管理", "field_name": "个人知识工厂", "sort_order": 4}, + {"trend_keyword": "数字生活", "field_name": "科技人文交叉", "sort_order": 5}, + {"trend_keyword": "科技人文", "field_name": "科技人文交叉", "sort_order": 6}, + {"trend_keyword": "个人成长", "field_name": "个人成长", "sort_order": 7}, + {"trend_keyword": "副业", "field_name": "个人成长", "sort_order": 8}, + {"trend_keyword": "AI创作", "field_name": "AI与效率", "sort_order": 9}, + {"trend_keyword": "未来工作", "field_name": "未来工作方式", "sort_order": 10}, + {"trend_keyword": "效率工具", "field_name": "AI与效率", "sort_order": 11}, + {"trend_keyword": "家庭教育", "field_name": "科技人文交叉", "sort_order": 12}, +] + +DEFAULT_CHINA_PAIN_TEMPLATES = { + "循环消费": "以旧换新流程繁琐、二手商品信任缺失、租赁市场不规范", + "低碳出行": "新能源车充电设施不足、城市规划不支持骑行、通勤距离长", + "干净饮食": "有机食品价格高、真伪难辨、外卖为主的生活方式难以改变", + "零浪费生活": "环保产品溢价高、可持续选择不便、漂绿营销难以分辨", + "绿色家电与节能": "绿色家电初期投入高、节能效果难量化、老旧小区改造难", + "碳普惠": "碳账户普及率低、减排量兑换吸引力不足、公众认知有限", + "环保科技产品": "绿色产品溢价68%难以承受、缺乏统一认证标准、担心漂绿", + "AI与效率": "AI工具选择困难、数据隐私担忧、学习成本高、实际效果难验证", +} + +DEFAULT_GLOBAL_RSS_KEYWORDS = ['sustainable', 'green', 'eco', 'circular', 'climate', 'carbon', 'zero waste', 'renewable', 'recycle', '环保', '可持续', '碳中和', '循环经济', '零浪费', '低碳', '生态'] + +DEFAULT_PRIORITY_WEIGHTS = {"audience_match": 0.3, "data_availability": 0.25, "uniqueness": 0.2, "executability": 0.15, "brand_fit": 0.1} + +DEFAULT_DOMAINS = ["远程工作", "AI工具", "可持续生活", "知识管理", "数字生活", "科技人文"] + +DEFAULT_SENSITIVE_WORDS = [ + {"word": "国家主席", "category": "political"}, + {"word": "政治局", "category": "political"}, + {"word": "常委", "category": "political"}, + {"word": "军委", "category": "political"}, + {"word": "统战部", "category": "political"}, + {"word": "颠覆国家", "category": "political"}, + {"word": "分裂主义", "category": "political"}, + {"word": "台独", "category": "political"}, + {"word": "疆独", "category": "political"}, + {"word": "藏独", "category": "political"}, + {"word": "赌博", "category": "prohibited"}, + {"word": "毒品", "category": "prohibited"}, + {"word": "迷药", "category": "prohibited"}, + {"word": "枪支", "category": "prohibited"}, + {"word": "炸药", "category": "prohibited"}, + {"word": "色情", "category": "prohibited"}, + {"word": "低俗", "category": "prohibited"}, + {"word": "反动", "category": "prohibited"}, + {"word": "邪教", "category": "prohibited"}, + {"word": "保证赚钱", "category": "misleading"}, + {"word": "一夜暴富", "category": "misleading"}, + {"word": "100%有效", "category": "misleading"}, + {"word": "包治百病", "category": "misleading"}, + {"word": "绝对正确", "category": "misleading"}, + {"word": "国家机密", "category": "legal"}, + {"word": "军事秘密", "category": "legal"}, + {"word": "绝密", "category": "legal"}, + {"word": "迷信", "category": "legal"}, +] + +DEFAULT_CONTENT_CLEAN_RULES = [ + {"rule_type": "thinking", "pattern": r"^(好的|好的,|好[的,]|我来|让我|我将|我这就).*?(?=\n|$)", "description": "AI思考模式1", "sort_order": 1}, + {"rule_type": "thinking", "pattern": r"^(以下|下面是|这是|为您|根据).*?(?=\n|$)", "description": "AI思考模式2", "sort_order": 2}, + {"rule_type": "preface", "pattern": r"^(基于|\u3010.*?\u3011|这里.*)", "description": "AI前缀模式", "sort_order": 3}, + {"rule_type": "verbosity", "pattern": r"^首先|^其次|^最后", "description": "AI废话-首先其次", "sort_order": 4}, + {"rule_type": "verbosity", "pattern": r"^总的来说$", "description": "AI废话-总的来说", "sort_order": 5}, + {"rule_type": "verbosity", "pattern": r"^值得注意的是$", "description": "AI废话-值得注意的是", "sort_order": 6}, + {"rule_type": "verbosity", "pattern": r"^换句话说$", "description": "AI废话-换句话说", "sort_order": 7}, + {"rule_type": "verbosity", "pattern": r"^总而言之$", "description": "AI废话-总而言之", "sort_order": 8}, + {"rule_type": "verbosity", "pattern": r"^简而言之$", "description": "AI废话-简而言之", "sort_order": 9}, + {"rule_type": "verbosity", "pattern": r"^一言以蔽之$", "description": "AI废话-一言以蔽之", "sort_order": 10}, + {"rule_type": "verbosity", "pattern": r"^可以说$", "description": "AI废话-可以说", "sort_order": 11}, + {"rule_type": "verbosity", "pattern": r"^不难发现$", "description": "AI废话-不难发现", "sort_order": 12}, + {"rule_type": "verbosity", "pattern": r"^由此可见$", "description": "AI废话-由此可见", "sort_order": 13}, + {"rule_type": "verbosity", "pattern": r"^综上所述$", "description": "AI废话-综上所述", "sort_order": 14}, + {"rule_type": "verbosity", "pattern": r"^通过以上", "description": "AI废话-通过以上", "sort_order": 15}, + {"rule_type": "html_thinking", "pattern": r"
]*>(好的|好的,|好[的,]|我来|让我|我将|我这就)", "description": "AI思考-HTML模式", "sort_order": 16}, +] + +class KeywordDomainMapResponse(BaseModel): + id: int + pattern: str + domain: str + sort_order: int + is_active: bool + created_at: Optional[datetime] = None + updated_at: Optional[datetime] = None + model_config = ConfigDict(from_attributes=True) + +class KeywordDomainMapCreate(BaseModel): + pattern: str + domain: str + sort_order: int = 0 + +class KeywordDomainMapUpdate(BaseModel): + pattern: Optional[str] = None + domain: Optional[str] = None + sort_order: Optional[int] = None + is_active: Optional[bool] = None + +class SensitiveWordResponse(BaseModel): + id: int + word: str + category: str + is_active: bool + added_by: Optional[str] = None + created_at: Optional[datetime] = None + model_config = ConfigDict(from_attributes=True) + +class SensitiveWordCreate(BaseModel): + word: str + category: str = "general" + +class ContentCleanRuleResponse(BaseModel): + id: int + rule_type: str + pattern: str + description: Optional[str] = None + is_active: bool + sort_order: int + created_at: Optional[datetime] = None + model_config = ConfigDict(from_attributes=True) + +class ContentCleanRuleCreate(BaseModel): + rule_type: str + pattern: str + description: Optional[str] = None + sort_order: int = 0 + +class ContentCleanRuleUpdate(BaseModel): + pattern: Optional[str] = None + description: Optional[str] = None + is_active: Optional[bool] = None + sort_order: Optional[int] = None + +class TrendFieldMappingResponse(BaseModel): + id: int + trend_keyword: str + field_name: str + sort_order: int + is_active: bool + created_at: Optional[datetime] = None + model_config = ConfigDict(from_attributes=True) + +class TrendFieldMappingCreate(BaseModel): + trend_keyword: str + field_name: str + sort_order: int = 0 + +class TrendFieldMappingUpdate(BaseModel): + trend_keyword: Optional[str] = None + field_name: Optional[str] = None + sort_order: Optional[int] = None + is_active: Optional[bool] = None + +class SystemConfigValueResponse(BaseModel): + key: str + value: Optional[str] = None + description: Optional[str] = None + +def _ensure_defaults(db: Session): + if db.query(KeywordDomainMap).count() == 0: + for item in DEFAULT_KEYWORD_DOMAIN_MAP: + db.add(KeywordDomainMap(**item)) + if db.query(SensitiveWord).count() == 0: + for item in DEFAULT_SENSITIVE_WORDS: + db.add(SensitiveWord(**item)) + if db.query(ContentCleanRule).count() == 0: + for item in DEFAULT_CONTENT_CLEAN_RULES: + db.add(ContentCleanRule(**item)) + if db.query(TrendFieldMapping).count() == 0: + for item in DEFAULT_TREND_FIELD_MAP: + db.add(TrendFieldMapping(**item)) + for name, pain in DEFAULT_CHINA_PAIN_TEMPLATES.items(): + cat = db.query(CollectorCategory).filter(CollectorCategory.name == name).first() + if cat and not cat.pain_template: + cat.pain_template = pain + db.commit() + +@router.get("/keyword-domain-map", response_model=List[KeywordDomainMapResponse]) +def list_keyword_domain_map(db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + _ensure_defaults(db) + return db.query(KeywordDomainMap).order_by(KeywordDomainMap.sort_order).all() + +@router.post("/keyword-domain-map", response_model=KeywordDomainMapResponse) +def create_keyword_domain_map(data: KeywordDomainMapCreate, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + item = KeywordDomainMap(**data.model_dump()) + db.add(item) + db.commit() + db.refresh(item) + return item + +@router.put("/keyword-domain-map/{item_id}", response_model=KeywordDomainMapResponse) +def update_keyword_domain_map(item_id: int, data: KeywordDomainMapUpdate, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + item = db.query(KeywordDomainMap).filter(KeywordDomainMap.id == item_id).first() + if not item: + raise HTTPException(status_code=404, detail="未找到") + for k, v in data.model_dump(exclude_unset=True).items(): + if v is not None: + setattr(item, k, v) + db.commit() + db.refresh(item) + return item + +@router.delete("/keyword-domain-map/{item_id}") +def delete_keyword_domain_map(item_id: int, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + item = db.query(KeywordDomainMap).filter(KeywordDomainMap.id == item_id).first() + if not item: + raise HTTPException(status_code=404, detail="未找到") + db.delete(item) + db.commit() + return {"message": "删除成功"} + +@router.get("/sensitive-words", response_model=List[SensitiveWordResponse]) +def list_sensitive_words(category: Optional[str] = None, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + _ensure_defaults(db) + q = db.query(SensitiveWord) + if category: + q = q.filter(SensitiveWord.category == category) + return q.order_by(SensitiveWord.id).all() + +@router.post("/sensitive-words", response_model=SensitiveWordResponse) +def create_sensitive_word(data: SensitiveWordCreate, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + item = SensitiveWord(**data.model_dump(), added_by=admin_user.username) + db.add(item) + db.commit() + db.refresh(item) + return item + +@router.put("/sensitive-words/{item_id}", response_model=SensitiveWordResponse) +def update_sensitive_word(item_id: int, enabled: bool = None, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + item = db.query(SensitiveWord).filter(SensitiveWord.id == item_id).first() + if not item: + raise HTTPException(status_code=404, detail="未找到") + if enabled is not None: + item.is_active = enabled + db.commit() + db.refresh(item) + return item + +@router.delete("/sensitive-words/{item_id}") +def delete_sensitive_word(item_id: int, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + item = db.query(SensitiveWord).filter(SensitiveWord.id == item_id).first() + if not item: + raise HTTPException(status_code=404, detail="未找到") + db.delete(item) + db.commit() + return {"message": "删除成功"} + +@router.get("/content-clean-rules", response_model=List[ContentCleanRuleResponse]) +def list_content_clean_rules(rule_type: Optional[str] = None, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + _ensure_defaults(db) + q = db.query(ContentCleanRule) + if rule_type: + q = q.filter(ContentCleanRule.rule_type == rule_type) + return q.order_by(ContentCleanRule.sort_order).all() + +@router.post("/content-clean-rules", response_model=ContentCleanRuleResponse) +def create_content_clean_rule(data: ContentCleanRuleCreate, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + item = ContentCleanRule(**data.model_dump()) + db.add(item) + db.commit() + db.refresh(item) + return item + +@router.put("/content-clean-rules/{item_id}", response_model=ContentCleanRuleResponse) +def update_content_clean_rule(item_id: int, data: ContentCleanRuleUpdate, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + item = db.query(ContentCleanRule).filter(ContentCleanRule.id == item_id).first() + if not item: + raise HTTPException(status_code=404, detail="未找到") + for k, v in data.model_dump(exclude_unset=True).items(): + if v is not None: + setattr(item, k, v) + db.commit() + db.refresh(item) + return item + +@router.delete("/content-clean-rules/{item_id}") +def delete_content_clean_rule(item_id: int, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + item = db.query(ContentCleanRule).filter(ContentCleanRule.id == item_id).first() + if not item: + raise HTTPException(status_code=404, detail="未找到") + db.delete(item) + db.commit() + return {"message": "删除成功"} + +@router.get("/trend-field-mapping", response_model=List[TrendFieldMappingResponse]) +def list_trend_field_mapping(db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + _ensure_defaults(db) + return db.query(TrendFieldMapping).filter(TrendFieldMapping.is_active == True).order_by(TrendFieldMapping.sort_order).all() + +@router.post("/trend-field-mapping", response_model=TrendFieldMappingResponse) +def create_trend_field_mapping(data: TrendFieldMappingCreate, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + item = TrendFieldMapping(**data.model_dump()) + db.add(item) + db.commit() + db.refresh(item) + return item + +@router.put("/trend-field-mapping/{item_id}", response_model=TrendFieldMappingResponse) +def update_trend_field_mapping(item_id: int, data: TrendFieldMappingUpdate, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + item = db.query(TrendFieldMapping).filter(TrendFieldMapping.id == item_id).first() + if not item: + raise HTTPException(status_code=404, detail="未找到") + for k, v in data.model_dump(exclude_unset=True).items(): + if v is not None: + setattr(item, k, v) + db.commit() + db.refresh(item) + return item + +@router.delete("/trend-field-mapping/{item_id}") +def delete_trend_field_mapping(item_id: int, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + item = db.query(TrendFieldMapping).filter(TrendFieldMapping.id == item_id).first() + if not item: + raise HTTPException(status_code=404, detail="未找到") + db.delete(item) + db.commit() + return {"message": "删除成功"} + +@router.get("/system-config/{key}", response_model=SystemConfigValueResponse) +def get_system_config(key: str, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + from ..models import SystemConfig + cfg = db.query(SystemConfig).filter(SystemConfig.key == key).first() + if not cfg: + default_map = { + "trend_domains": json.dumps(DEFAULT_DOMAINS), + "rss_default_keywords": json.dumps(DEFAULT_GLOBAL_RSS_KEYWORDS), + "priority_weights": json.dumps(DEFAULT_PRIORITY_WEIGHTS), + } + if key in default_map: + return SystemConfigValueResponse(key=key, value=default_map[key], description=f"系统默认配置 - {key}") + raise HTTPException(status_code=404, detail="未找到") + return SystemConfigValueResponse(key=cfg.key, value=cfg.value, description=cfg.description) + +@router.put("/system-config/{key}") +def update_system_config(key: str, value: str, description: Optional[str] = None, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + from ..models import SystemConfig + cfg = db.query(SystemConfig).filter(SystemConfig.key == key).first() + if cfg: + cfg.value = value + if description is not None: + cfg.description = description + else: + cfg = SystemConfig(key=key, value=value, description=description or key) + db.add(cfg) + db.commit() + return {"message": "保存成功", "key": key, "value": value} + +@router.get("/system-configs") +def list_system_configs(db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + from ..models import SystemConfig + configs = db.query(SystemConfig).order_by(SystemConfig.key).all() + defaults = { + "trend_domains": json.dumps(DEFAULT_DOMAINS), + "rss_default_keywords": json.dumps(DEFAULT_GLOBAL_RSS_KEYWORDS), + "priority_weights": json.dumps(DEFAULT_PRIORITY_WEIGHTS), + } + result = [] + for c in configs: + result.append({"key": c.key, "value": c.value, "description": c.description, "is_default": False}) + for k, v in defaults.items(): + if not any(x["key"] == k for x in result): + result.append({"key": k, "value": v, "description": f"系统默认 - {k}", "is_default": True}) + return result \ No newline at end of file diff --git a/platform/backend/app/api/prompt_configs.py b/platform/backend/app/api/prompt_configs.py new file mode 100644 index 0000000..72bd39e --- /dev/null +++ b/platform/backend/app/api/prompt_configs.py @@ -0,0 +1,117 @@ +from fastapi import APIRouter, Depends, HTTPException +from sqlalchemy.orm import Session +from typing import List, Optional +from pydantic import BaseModel, ConfigDict +from datetime import datetime + +from ..database import get_db +from ..models import PromptConfig +from .auth import get_current_admin + +router = APIRouter(prefix="/api/admin/prompt-configs", tags=["admin"]) + + +class PromptConfigBase(BaseModel): + key: str + module_id: Optional[str] = None + category: str = "prompt" + version: str = "v1" + content: str + variables: List[dict] = [] + description: Optional[str] = None + enabled: bool = True + temperature: Optional[float] = None + max_tokens: Optional[int] = None + created_by: Optional[str] = None + + +class PromptConfigCreate(PromptConfigBase): + pass + + +class PromptConfigUpdate(BaseModel): + content: Optional[str] = None + variables: Optional[List[dict]] = None + description: Optional[str] = None + enabled: Optional[bool] = None + temperature: Optional[float] = None + max_tokens: Optional[int] = None + + +class PromptConfigResponse(PromptConfigBase): + id: int + created_at: Optional[datetime] = None + updated_at: Optional[datetime] = None + + model_config = ConfigDict(from_attributes=True) + + +@router.get("", response_model=List[PromptConfigResponse]) +def list_prompts( + module_id: Optional[str] = None, + category: Optional[str] = None, + db: Session = Depends(get_db), + admin_user=Depends(get_current_admin), +): + q = db.query(PromptConfig) + if module_id: + q = q.filter(PromptConfig.module_id.in_([module_id, "all"])) + if category: + q = q.filter(PromptConfig.category == category) + prompts = q.order_by(PromptConfig.module_id, PromptConfig.key).all() + if not prompts: + _ensure_defaults(db) + prompts = q.order_by(PromptConfig.module_id, PromptConfig.key).all() + return [PromptConfigResponse.model_validate(p) for p in prompts] + + +@router.get("/{prompt_id}", response_model=PromptConfigResponse) +def get_prompt(prompt_id: int, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + prompt = db.query(PromptConfig).filter(PromptConfig.id == prompt_id).first() + if not prompt: + raise HTTPException(status_code=404, detail="未找到该配置") + return prompt + + +@router.post("", response_model=PromptConfigResponse) +def create_prompt(data: PromptConfigCreate, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + existing = db.query(PromptConfig).filter(PromptConfig.key == data.key).first() + if existing: + raise HTTPException(status_code=400, detail=f"key '{data.key}' 已存在") + prompt = PromptConfig(**data.model_dump()) + db.add(prompt) + db.commit() + db.refresh(prompt) + return prompt + + +@router.put("/{prompt_id}", response_model=PromptConfigResponse) +def update_prompt(prompt_id: int, data: PromptConfigUpdate, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + prompt = db.query(PromptConfig).filter(PromptConfig.id == prompt_id).first() + if not prompt: + raise HTTPException(status_code=404, detail="未找到该配置") + for field, value in data.model_dump(exclude_unset=True).items(): + if value is not None: + setattr(prompt, field, value) + db.commit() + db.refresh(prompt) + return prompt + + +@router.delete("/{prompt_id}") +def delete_prompt(prompt_id: int, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + prompt = db.query(PromptConfig).filter(PromptConfig.id == prompt_id).first() + if not prompt: + raise HTTPException(status_code=404, detail="未找到该配置") + db.delete(prompt) + db.commit() + return {"message": "删除成功"} + + +def _ensure_defaults(db: Session): + existing_keys = {p.key for p in db.query(PromptConfig).all()} + from .task_configs import DEFAULT_PROMPTS + for p in DEFAULT_PROMPTS: + if p["key"] not in existing_keys: + db.add(PromptConfig(**p)) + db.commit() \ No newline at end of file diff --git a/platform/backend/app/api/system.py b/platform/backend/app/api/system.py index e0cbc03..7617459 100644 --- a/platform/backend/app/api/system.py +++ b/platform/backend/app/api/system.py @@ -9,7 +9,7 @@ from typing import Dict, Any, List, Optional import os import json from ..database import get_db -from ..models import Topic, Article +from ..models import Topic, Article, TaskConfig, TaskLog from ..core.generator import run_creator, get_generator_status from ..core.optimizer import run_optimizer, get_optimizer_status from ..core.collector import run_collector, get_collector_status @@ -61,7 +61,7 @@ def get_status(db: Session = Depends(get_db)): } @router.post("/generate/run", dependencies=[Depends(get_current_user)]) -def trigger_generation(topic_id: str = Body(None, embed=True), db: Session = Depends(get_db), current_user=Depends(get_current_user)): +def trigger_generation(topic_id: Optional[str] = None, db: Session = Depends(get_db), current_user=Depends(get_current_user)): logger.info(f"Generation triggered by {current_user.username}, topic_id={topic_id}") try: result = run_creator(topic_id) @@ -93,7 +93,7 @@ def collection_status(): return status @router.post("/review/run", dependencies=[Depends(get_current_user)]) -def trigger_review(topic_ids: List[str] = Body(None, embed=True), db: Session = Depends(get_db), current_user=Depends(get_current_user)): +def trigger_review(topic_ids: Optional[List[str]] = None, db: Session = Depends(get_db), current_user=Depends(get_current_user)): try: result = run_optimizer(topic_ids) return {"message": "合规审查已后台启动", "pid": result.get("pid")} @@ -247,47 +247,61 @@ def get_scheduler_status(): @router.get("/modules/status", dependencies=[Depends(get_current_user)]) -def get_modules_status(): - today_str = date.today().isoformat() - log_based: dict = { - "scheduled_collect": {"name": "📡 内容采集", "log": LOGS_DIR / f"collector_{today_str}.log"}, - "scheduled_refresh_search_cache": {"name": "🔍 搜索缓存", "log": LOGS_DIR / f"opencode_search_{today_str}.log"}, - "scheduled_fetch_trends": {"name": "🔥 热点趋势", "log": LOGS_DIR / f"trends_{today_str}.log"}, - "scheduled_generate": {"name": "🤖 内容创作", "log": LOGS_DIR / f"creator_{today_str}.log"}, - "scheduled_optimize": {"name": "🔍 合规审查", "log": LOGS_DIR / f"optimizer_{today_str}.log"}, - "scheduled_optimize_sources": {"name": "📡 信息源优化", "log": LOGS_DIR / f"optimizer_sources_{today_str}.log"}, - "scheduled_metrics_sync": {"name": "📊 指标同步", "log": LOGS_DIR / f"metrics_sync_{today_str}.log"}, +def get_modules_status(db: Session = Depends(get_db)): + configs = db.query(TaskConfig).all() + config_map = {c.module_id: c for c in configs} + + MODULE_META = { + "scheduled_refresh_search_cache": {"name": "🔍 搜索缓存", "cron": "01:00", "params_desc": {"refresh_queries": "搜索关键词列表"}}, + "scheduled_fetch_trends": {"name": "🔥 热点趋势", "cron": "01:10", "params_desc": {}}, + "scheduled_collect": {"name": "📡 内容采集", "cron": "01:30", "params_desc": {"max_topics": "最大选题数", "categories": "采集类别"}}, + "scheduled_generate": {"name": "🤖 内容创作", "cron": "02:00", "params_desc": {"auto_review": "自动合规审查"}}, + "scheduled_optimize": {"name": "🔍 合规审查", "cron": "03:00", "params_desc": {"auto_pass_threshold": "自动通过分数阈值"}}, + "scheduled_optimize_sources": {"name": "📡 信息源优化", "cron": "05:00", "params_desc": {}}, + "scheduled_metrics_sync": {"name": "📊 指标同步", "cron": "06:00", "params_desc": {}}, } - jobs = {j['id']: j for j in scheduler.get_jobs()} + modules = [] - for mod_id, cfg in log_based.items(): - log_file = cfg["log"] - last_run = None - task_count = 0 - success_rate = None - if log_file.exists(): - mtime = datetime.fromtimestamp(log_file.stat().st_mtime) - last_run = mtime.strftime("%Y-%m-%d %H:%M") - content = log_file.read_text(encoding="utf-8", errors="ignore") - task_count = content.count("完成") + content.count("success") + content.count("SUCCESS") - total = task_count + content.count("失败") + content.count("failed") + content.count("ERROR") - success_rate = round(task_count / total * 100) if total > 0 else None - status = "running" if mod_id in jobs else "stopped" - job = jobs.get(mod_id) - next_run = None - if job and job.get("next_run_time"): - try: - next_dt = datetime.fromisoformat(job["next_run_time"]) - next_run = next_dt.strftime("%Y-%m-%d %H:%M") - except Exception: - next_run = job["next_run_time"] + for mod_id, meta in MODULE_META.items(): + cfg = config_map.get(mod_id) + latest = db.query(TaskLog).filter(TaskLog.module_id == mod_id).order_by(TaskLog.started_at.desc()).first() + next_run = _get_next_run(mod_id) + + total = db.query(TaskLog).filter(TaskLog.module_id == mod_id).count() + success = db.query(TaskLog).filter(TaskLog.module_id == mod_id, TaskLog.status == "success").count() + failed = db.query(TaskLog).filter(TaskLog.module_id == mod_id, TaskLog.status == "failed").count() + running = db.query(TaskLog).filter(TaskLog.module_id == mod_id, TaskLog.status == "running").count() + modules.append({ - "id": mod_id, - "title": cfg["name"], - "status": status, - "last_run": last_run or "从未运行", - "next_run": next_run or "待计划", - "task_count": task_count, - "success_rate": success_rate if success_rate is not None else 0, + "module_id": mod_id, + "title": meta["name"], + "enabled": cfg.enabled if cfg else True, + "params": cfg.params if cfg else {}, + "params_desc": meta["params_desc"], + "schedule": cfg.schedule if cfg else meta["cron"], + "cron_default": meta["cron"], + "status": "running" if running else ("stopped" if not (cfg and cfg.enabled) else "idle"), + "last_run": latest.started_at.strftime("%Y-%m-%d %H:%M") if latest and latest.started_at else None, + "last_status": latest.status if latest else None, + "last_message": latest.message if latest else None, + "last_result": latest.result_data if latest else None, + "next_run": next_run, + "total_runs": total, + "success_runs": success, + "failed_runs": failed, + "running": running, }) - return {"modules": modules, "scheduler": {"running": scheduler._started, "jobs": scheduler.get_jobs()}} + + jobs = scheduler.get_jobs() + return {"modules": modules, "scheduler": {"running": scheduler._started, "jobs": jobs}} + + +def _get_next_run(mod_id: str) -> Optional[str]: + for job in scheduler.get_jobs(): + if job["id"] == mod_id and job["next_run_time"]: + try: + dt = datetime.fromisoformat(job["next_run_time"]) + return dt.strftime("%Y-%m-%d %H:%M") + except Exception: + return job["next_run_time"] + return None diff --git a/platform/backend/app/api/task_configs.py b/platform/backend/app/api/task_configs.py new file mode 100644 index 0000000..bfeb005 --- /dev/null +++ b/platform/backend/app/api/task_configs.py @@ -0,0 +1,66 @@ +from fastapi import APIRouter, Depends, HTTPException +from sqlalchemy.orm import Session +from typing import List, Optional + +from ..database import get_db +from ..models import TaskConfig, TaskLog +from ..schemas import TaskConfigBase, TaskConfigUpdate, TaskConfigResponse, TaskLogResponse +from .auth import get_current_admin + +router = APIRouter(prefix="/api/admin/task-configs", tags=["admin"]) + +DEFAULT_CONFIGS = { + "scheduled_refresh_search_cache": {"name": "🔍 搜索缓存", "cron": "01:00", "params": {}}, + "scheduled_fetch_trends": {"name": "🔥 热点趋势", "cron": "01:10", "params": {}}, + "scheduled_collect": {"name": "📡 内容采集", "cron": "01:30", "params": {"max_topics": 20}}, + "scheduled_generate": {"name": "🤖 内容创作", "cron": "02:00", "params": {"auto_review": True}}, + "scheduled_optimize": {"name": "🔍 合规审查", "cron": "03:00", "params": {"auto_pass_threshold": 80}}, + "scheduled_optimize_sources": {"name": "📡 信息源优化", "cron": "05:00", "params": {}}, + "scheduled_metrics_sync": {"name": "📊 指标同步", "cron": "06:00", "params": {}}, +} + +@router.get("", response_model=List[TaskConfigResponse]) +def list_configs(db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + configs = db.query(TaskConfig).order_by(TaskConfig.id).all() + if not configs: + _ensure_defaults(db) + configs = db.query(TaskConfig).order_by(TaskConfig.id).all() + return [TaskConfigResponse.model_validate(c) for c in configs] + +@router.get("/{module_id}", response_model=TaskConfigResponse) +def get_config(module_id: str, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + cfg = db.query(TaskConfig).filter(TaskConfig.module_id == module_id).first() + if not cfg: + _ensure_defaults(db) + cfg = db.query(TaskConfig).filter(TaskConfig.module_id == module_id).first() + return cfg + +@router.put("/{module_id}", response_model=TaskConfigResponse) +def update_config(module_id: str, data: TaskConfigUpdate, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + cfg = db.query(TaskConfig).filter(TaskConfig.module_id == module_id).first() + if not cfg: + _ensure_defaults(db) + cfg = db.query(TaskConfig).filter(TaskConfig.module_id == module_id).first() + if data.enabled is not None: + cfg.enabled = data.enabled + if data.params is not None: + cfg.params = data.params + if data.schedule is not None: + cfg.schedule = data.schedule + if data.last_modified_by: + cfg.last_modified_by = data.last_modified_by + db.commit() + db.refresh(cfg) + return cfg + +@router.get("/history/{module_id}", response_model=List[TaskLogResponse]) +def get_module_history(module_id: str, db: Session = Depends(get_db), admin_user=Depends(get_current_admin), limit: int = 20): + logs = db.query(TaskLog).filter(TaskLog.module_id == module_id).order_by(TaskLog.started_at.desc()).limit(limit).all() + return [TaskLogResponse.model_validate(l) for l in logs] + +def _ensure_defaults(db: Session): + existing = {c.module_id for c in db.query(TaskConfig).all()} + for mid, info in DEFAULT_CONFIGS.items(): + if mid not in existing: + db.add(TaskConfig(module_id=mid, enabled=True, params=info.get("params", {}), schedule=info.get("cron", ""))) + db.commit() \ No newline at end of file diff --git a/platform/backend/app/api/task_logs.py b/platform/backend/app/api/task_logs.py index 7dacd58..d211791 100644 --- a/platform/backend/app/api/task_logs.py +++ b/platform/backend/app/api/task_logs.py @@ -1,55 +1,82 @@ -from fastapi import APIRouter, Depends, HTTPException, Request +from fastapi import APIRouter, Depends, HTTPException from sqlalchemy.orm import Session +from sqlalchemy import or_ from typing import List, Optional +from datetime import datetime, timezone from ..database import get_db from ..models import TaskLog from ..schemas import TaskLogBase, TaskLogResponse from .auth import get_current_admin -router = APIRouter(prefix="/api/admin/tasklogs", tags=["admin"]) +router = APIRouter(prefix="/api/admin/task-logs", tags=["admin"]) + +MODULES = { + "scheduled_refresh_search_cache": "🔍 搜索缓存", + "scheduled_fetch_trends": "🔥 热点趋势", + "scheduled_collect": "📡 内容采集", + "scheduled_generate": "🤖 内容创作", + "scheduled_optimize": "🔍 合规审查", + "scheduled_optimize_sources": "📡 信息源优化", + "scheduled_metrics_sync": "📊 指标同步", +} @router.get("", response_model=List[TaskLogResponse]) def list_task_logs( - request: Request, db: Session = Depends(get_db), - admin_user = Depends(get_current_admin), - topic_id: Optional[str] = None, - task_name: Optional[str] = None, - status: Optional[str] = None + admin_user=Depends(get_current_admin), + module_id: Optional[str] = None, + status: Optional[str] = None, + date: Optional[str] = None, + limit: int = 50, ): - """获取任务日志列表(可过滤)""" query = db.query(TaskLog) - if topic_id: - query = query.filter(TaskLog.topic_id == topic_id) - if task_name: - query = query.filter(TaskLog.task_name == task_name) + if module_id: + query = query.filter(TaskLog.module_id == module_id) if status: query = query.filter(TaskLog.status == status) - logs = query.order_by(TaskLog.started_at.desc()).all() + if date: + try: + dt = datetime.strptime(date, "%Y-%m-%d").replace(tzinfo=timezone.utc) + next_day = datetime(dt.year, dt.month, dt.day + 1, tzinfo=timezone.utc) + query = query.filter(TaskLog.started_at >= dt, TaskLog.started_at < next_day) + except ValueError: + pass + logs = query.order_by(TaskLog.started_at.desc()).limit(limit).all() return [TaskLogResponse.model_validate(l) for l in logs] +@router.get("/modules") +def list_modules(db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): + today = datetime.now(timezone.utc).date().isoformat() + result = [] + for mid, name in MODULES.items(): + latest = db.query(TaskLog).filter(TaskLog.module_id == mid).order_by(TaskLog.started_at.desc()).first() + total = db.query(TaskLog).filter(TaskLog.module_id == mid).count() + success = db.query(TaskLog).filter(TaskLog.module_id == mid, TaskLog.status == "success").count() + failed = db.query(TaskLog).filter(TaskLog.module_id == mid, TaskLog.status == "failed").count() + running = db.query(TaskLog).filter(TaskLog.module_id == mid, TaskLog.status == "running").count() + result.append({ + "module_id": mid, + "name": name, + "last_run": latest.started_at.isoformat() if latest and latest.started_at else None, + "last_status": latest.status if latest else None, + "last_message": latest.message if latest else None, + "total_runs": total, + "success_runs": success, + "failed_runs": failed, + "running": running, + }) + return result + @router.get("/{log_id}", response_model=TaskLogResponse) -def get_task_log( - log_id: int, - request: Request, - db: Session = Depends(get_db), - admin_user = Depends(get_current_admin) -): - """获取单个任务日志详情""" +def get_task_log(log_id: int, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): log = db.query(TaskLog).filter(TaskLog.id == log_id).first() if not log: - raise HTTPException(status_code=404, detail="日志不存在") + raise HTTPException(status_code=404, detail="记录不存在") return log @router.post("", response_model=TaskLogResponse) -def create_task_log( - log_data: TaskLogBase, - request: Request, - db: Session = Depends(get_db), - admin_user = Depends(get_current_admin) -): - """创建任务日志(用于手动记录)""" +def create_task_log(log_data: TaskLogBase, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): log = TaskLog(**log_data.model_dump()) db.add(log) db.commit() @@ -57,35 +84,22 @@ def create_task_log( return log @router.put("/{log_id}", response_model=TaskLogResponse) -def update_task_log( - log_id: int, - log_update: TaskLogBase, - request: Request, - db: Session = Depends(get_db), - admin_user = Depends(get_current_admin) -): - """更新任务日志""" +def update_task_log(log_id: int, log_update: TaskLogBase, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): log = db.query(TaskLog).filter(TaskLog.id == log_id).first() if not log: - raise HTTPException(status_code=404, detail="日志不存在") - update_data = log_update.model_dump(exclude_unset=True) - for field, value in update_data.items(): + raise HTTPException(status_code=404, detail="记录不存在") + data = log_update.model_dump(exclude_unset=True) + for field, value in data.items(): setattr(log, field, value) db.commit() db.refresh(log) return log @router.delete("/{log_id}") -def delete_task_log( - log_id: int, - request: Request, - db: Session = Depends(get_db), - admin_user = Depends(get_current_admin) -): - """删除任务日志""" +def delete_task_log(log_id: int, db: Session = Depends(get_db), admin_user=Depends(get_current_admin)): log = db.query(TaskLog).filter(TaskLog.id == log_id).first() if not log: - raise HTTPException(status_code=404, detail="日志不存在") + raise HTTPException(status_code=404, detail="记录不存在") db.delete(log) db.commit() - return {"message": "删除成功"} + return {"message": "删除成功"} \ No newline at end of file diff --git a/platform/backend/app/api/tasks.py b/platform/backend/app/api/tasks.py index 26ed53f..788192b 100644 --- a/platform/backend/app/api/tasks.py +++ b/platform/backend/app/api/tasks.py @@ -213,6 +213,25 @@ def get_module_detail(module_id: str, db: Session = Depends(get_db), current_use LOGS_DIR = ROOT / "automation" / "logs" today_str = date_mod.today().isoformat() + try: + return _get_module_detail_data(module_id, db, ROOT, DATA_DIR, LOGS_DIR, today_str) + except Exception as e: + return { + "module_id": module_id, + "title": module_id, + "description": "", + "status": "stopped", + "inputs": {}, + "outputs": {"error": str(e)}, + "history": [], + "log_excerpt": "", + } + +def _get_module_detail_data(module_id: str, db, ROOT, DATA_DIR, LOGS_DIR, today_str): + from datetime import datetime as dt_mod, date as date_mod + from pathlib import Path as PathMod + import json as json_mod + import re as re_mod MODULE_META = { "scheduled_refresh_search_cache": {"name": "🔍 搜索缓存", "description": "通过 opencode webfetch 联网搜索,刷新 8 个分类的搜索缓存,供内容采集器使用"}, "scheduled_fetch_trends": {"name": "🔥 热点趋势", "description": "从百度、微博、知乎实时热搜 API 抓取当天热点,LLM 补充,存入 trends.json"}, @@ -279,7 +298,7 @@ def get_module_detail(module_id: str, db: Session = Depends(get_db), current_use pass elif module_id == "scheduled_collect": - from ..models import Topic, CollectorCategory + from ..models import CollectorCategory pending = db.query(Topic).filter(Topic.status.in_(["pending", "待处理"])).count() total_topics = db.query(Topic).count() cats = db.query(CollectorCategory).filter(CollectorCategory.is_active == True).all() @@ -296,16 +315,19 @@ def get_module_detail(module_id: str, db: Session = Depends(get_db), current_use pending_t = db.query(Topic).filter(Topic.status.in_(["pending", "待处理"])).count() inputs["待创作选题"] = pending_t outputs["待审查"] = review - from ..models import Article - recent_articles = db.query(Article, Topic.title.label("topic_title")).join(Topic, Article.topic_id == Topic.id, isouter=True).order_by(Article.created_at.desc()).limit(5).all() - outputs["最新文章"] = [] - seen_articles = set() - for a in recent_articles: - art = a.Article if hasattr(a, 'Article') else a[0] - tid = a.topic_title if hasattr(a, 'topic_title') else (a[1] if len(a) > 1 else "") - if art.id not in seen_articles: - seen_articles.add(art.id) - outputs["最新文章"].append({"id": art.id, "platform": art.platform, "topic": tid, "status": art.status, "created": art.created_at.isoformat() if art.created_at else ""}) + try: + from ..models import Article + recent_articles = db.query(Article, Topic.title.label("topic_title")).join(Topic, Article.topic_id == Topic.id, isouter=True).order_by(Article.created_at.desc()).limit(5).all() + outputs["最新文章"] = [] + seen_articles = set() + for a in recent_articles: + art = a[0] + tid = a[1] if len(a) > 1 else "" + if art.id not in seen_articles: + seen_articles.add(art.id) + outputs["最新文章"].append({"id": art.id, "platform": art.platform, "topic": tid, "status": art.status, "created": art.created_at.isoformat() if art.created_at else ""}) + except Exception as e: + outputs["最新文章_错误"] = str(e) elif module_id == "scheduled_optimize": outputs["待审查选题"] = db.query(Topic).filter(Topic.status.in_(["review", "待审查"])).count() diff --git a/platform/backend/app/core/prompt_loader.py b/platform/backend/app/core/prompt_loader.py new file mode 100644 index 0000000..34e93e8 --- /dev/null +++ b/platform/backend/app/core/prompt_loader.py @@ -0,0 +1,108 @@ +import os +from typing import Optional, Dict, Any + +USE_POSTGRES = os.getenv('USE_POSTGRES', 'true').lower() == 'true' +if USE_POSTGRES: + from sqlalchemy import create_engine + from sqlalchemy.orm import sessionmaker + POSTGRES_CONFIG = { + 'host': os.getenv('PG_HOST', '127.0.0.1'), + 'port': os.getenv('PG_PORT', '5432'), + 'database': os.getenv('PG_DATABASE', 'yzr_nr'), + 'user': os.getenv('PG_USER', 'yzr_nr'), + 'password': os.getenv('PG_PASSWORD', 'aTX3WKKnPfRnM5PC') + } + SQLALCHEMY_DATABASE_URL = ( + f"postgresql://{POSTGRES_CONFIG['user']}:{POSTGRES_CONFIG['password']}" + f"@{POSTGRES_CONFIG['host']}:{POSTGRES_CONFIG['port']}/{POSTGRES_CONFIG['database']}" + ) + _engine = create_engine(SQLALCHEMY_DATABASE_URL, pool_pre_ping=True) +else: + from pathlib import Path + PROJECT_ROOT = Path(__file__).resolve().parent.parent.parent.parent + DATA_DIR = os.getenv('DATA_DIR', str(PROJECT_ROOT / 'data')) + os.makedirs(DATA_DIR, exist_ok=True) + DB_PATH = os.path.join(DATA_DIR, 'yzr.db') + from sqlalchemy import create_engine + _engine = create_engine(f"sqlite:///{DB_PATH}", connect_args={"check_same_thread": False}) + +_SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=_engine) + + +def _get_session(): + session = _SessionLocal() + try: + return session + except: + session.close() + raise + + +_PROMPT_CACHE: Dict[str, Dict[str, Any]] = {} +_CACHE_VERSION = "v2" + + +def load_prompt_config(key: str, module_id: Optional[str] = None) -> Optional[Dict[str, Any]]: + if key in _PROMPT_CACHE and _PROMPT_CACHE[key].get("_v") == _CACHE_VERSION: + return _PROMPT_CACHE[key] + try: + session = _get_session() + try: + from .models import PromptConfig + prompt = session.query(PromptConfig).filter( + PromptConfig.key == key, + PromptConfig.enabled == True + ).first() + if not prompt and module_id: + prompt = session.query(PromptConfig).filter( + PromptConfig.key == key, + PromptConfig.module_id.in_([module_id, "all"]), + PromptConfig.enabled == True + ).first() + if prompt: + result = { + "_v": _CACHE_VERSION, + "content": prompt.content, + "temperature": prompt.temperature, + "max_tokens": prompt.max_tokens, + "variables": prompt.variables or [], + "description": prompt.description, + "category": prompt.category, + } + _PROMPT_CACHE[key] = result + return result + finally: + session.close() + except Exception as e: + import warnings + warnings.warn(f"load_prompt_config({key}) failed: {e}") + return None + + +def clear_prompt_cache(): + _PROMPT_CACHE.clear() + + +def get_prompt(key: str, module_id: Optional[str] = None, **kwargs) -> str: + cfg = load_prompt_config(key, module_id) + if cfg: + content = cfg["content"] + for var in cfg.get("variables", []): + name = var.get("name") + if name and name in kwargs: + content = content.replace("{" + name + "}", str(kwargs[name])) + elif name: + default = var.get("default_value", "") + content = content.replace("{" + name + "}", str(default)) + return content + return "" + + +def get_llm_params(key: str, module_id: Optional[str] = None) -> Dict[str, Any]: + cfg = load_prompt_config(key, module_id) + if cfg: + return { + "temperature": cfg.get("temperature"), + "max_tokens": cfg.get("max_tokens"), + } + return {} \ No newline at end of file diff --git a/platform/backend/app/core/scheduler.py b/platform/backend/app/core/scheduler.py index 1f03c61..1915358 100644 --- a/platform/backend/app/core/scheduler.py +++ b/platform/backend/app/core/scheduler.py @@ -2,11 +2,15 @@ 定时任务调度器 基于 APScheduler,支持在 FastAPI 生命周期内运行定时任务 """ -import os -import sys -import logging -from datetime import datetime +import os, sys, logging, json from pathlib import Path +from datetime import datetime, timezone + +PROJECT_ROOT = Path(__file__).parent.parent.parent.parent.parent +sys.path.insert(0, str(PROJECT_ROOT / 'scripts')) +sys.path.insert(0, str(PROJECT_ROOT)) +from prompt_loader import get_prompt, get_prompt_params + from apscheduler.schedulers.background import BackgroundScheduler from apscheduler.triggers.cron import CronTrigger from .generator import run_creator_blocking @@ -15,6 +19,69 @@ from .collector import run_collector_blocking logger = logging.getLogger(__name__) +MODULES = { + "scheduled_refresh_search_cache": {"name": "🔍 搜索缓存", "cron": "01:00"}, + "scheduled_fetch_trends": {"name": "🔥 热点趋势", "cron": "01:10"}, + "scheduled_collect": {"name": "📡 内容采集", "cron": "01:30"}, + "scheduled_generate": {"name": "🤖 内容创作", "cron": "02:00"}, + "scheduled_optimize": {"name": "🔍 合规审查", "cron": "03:00"}, + "scheduled_optimize_sources": {"name": "📡 信息源优化", "cron": "05:00"}, + "scheduled_metrics_sync": {"name": "📊 指标同步", "cron": "06:00"}, +} + +def _log_task(module_id: str, status: str, message: str = None, + error_trace: str = None, result_data: dict = None, + started_at: datetime = None, finished_at: datetime = None, + triggered_by: str = "scheduler", next_run_time: datetime = None): + try: + from ..database import SessionLocal + from ..models import TaskLog + db = SessionLocal() + try: + duration = None + if started_at and finished_at: + duration = int((finished_at - started_at).total_seconds()) + log = TaskLog( + module_id=module_id, + task_name=MODULES.get(module_id, {}).get("name", module_id), + status=status, + message=message, + error_trace=error_trace, + triggered_by=triggered_by, + result_data=result_data or {}, + started_at=started_at or datetime.now(timezone.utc), + finished_at=finished_at, + duration=duration, + next_run_time=next_run_time, + ) + db.add(log) + db.commit() + finally: + db.close() + except Exception: + pass + +def _wrap_task(module_id: str, target_fn, *args, **kwargs): + started = datetime.now(timezone.utc) + status = "running" + error_trace = None + result_data = None + try: + result = target_fn(*args, **kwargs) + status = "success" + if isinstance(result, dict): + result_data = {k: v for k, v in result.items() if isinstance(v, (str, int, float, bool, list, dict)) and k not in ("stdout", "stderr")} + return result + except Exception as e: + status = "failed" + import traceback + error_trace = traceback.format_exc() + raise + finally: + _log_task(module_id, status=status, message=None, error_trace=error_trace, + result_data=result_data, started_at=started, + finished_at=datetime.now(timezone.utc)) + class TaskScheduler: def __init__(self): self.scheduler = BackgroundScheduler() @@ -24,64 +91,47 @@ class TaskScheduler: if self._started: logger.warning("Scheduler already started") return - # 使用 CronTrigger 设置每日固定时间点 - # 顺序: 搜索缓存(01:00)→趋势(01:10)→采集(01:30)→创作(02:00)→审查(03:00)→源优化(05:00)→指标(06:00) - self.scheduler.add_job( - self._run_refresh_search_cache, - CronTrigger(hour=1, minute=0), - id='scheduled_refresh_search_cache', - replace_existing=True, - max_instances=1, - coalesce=True - ) - self.scheduler.add_job( - self._run_fetch_trends, - CronTrigger(hour=1, minute=10), - id='scheduled_fetch_trends', - replace_existing=True, - max_instances=1, - coalesce=True - ) - self.scheduler.add_job( - self._run_collect, - CronTrigger(hour=1, minute=30), - id='scheduled_collect', - ) - self.scheduler.add_job( - self._run_generate, - CronTrigger(hour=2, minute=0), - id='scheduled_generate', - replace_existing=True, - max_instances=1, - coalesce=True - ) - self.scheduler.add_job( - self._run_optimize, - CronTrigger(hour=3, minute=0), - id='scheduled_optimize', - replace_existing=True, - max_instances=1, - coalesce=True - ) - self.scheduler.add_job( - self._run_optimize_sources, - CronTrigger(hour=5, minute=0), - id='scheduled_optimize_sources', - replace_existing=True, - max_instances=1, - coalesce=True - ) - self.scheduler.add_job( - self._run_metrics_sync, - CronTrigger(hour=6, minute=0), - id='scheduled_metrics_sync', - replace_existing=True, - max_instances=1, - coalesce=True - ) + from ..database import SessionLocal + from ..models import TaskConfig + db = SessionLocal() + try: + configs = {c.module_id: c for c in db.query(TaskConfig).all()} + finally: + db.close() + + MODULE_JOBS = [ + ("scheduled_refresh_search_cache", self._run_refresh_search_cache, "搜索缓存"), + ("scheduled_fetch_trends", self._run_fetch_trends, "热点趋势"), + ("scheduled_collect", self._run_collect, "内容采集"), + ("scheduled_generate", self._run_generate, "内容创作"), + ("scheduled_optimize", self._run_optimize, "合规审查"), + ("scheduled_optimize_sources", self._run_optimize_sources, "信息源优化"), + ("scheduled_metrics_sync", self._run_metrics_sync, "指标同步"), + ] + + for module_id, fn, name in MODULE_JOBS: + cfg = configs.get(module_id) + if cfg and not cfg.enabled: + logger.info(f"跳过禁用任务: {module_id}") + continue + schedule = (cfg.schedule if cfg else None) or MODULES.get(module_id, {}).get("cron", "01:00") + try: + hour, minute = map(int, schedule.split(":")) + except (ValueError, AttributeError): + hour, minute = 1, 0 + self.scheduler.add_job( + fn, + CronTrigger(hour=hour, minute=minute), + id=module_id, + replace_existing=True, + max_instances=1, + coalesce=True + ) + logger.info(f"调度任务: {module_id} -> {schedule}") + self.scheduler.start() self._started = True - logger.info("Scheduler started: 01:00 search 01:10 trends 01:30 collect 02:00 create 03:00 review 05:00 sources 06:00 metrics") + logger.info("Scheduler started with dynamic schedule from TaskConfig") def shutdown(self): if self.scheduler.running: self.scheduler.shutdown() @@ -89,6 +139,8 @@ class TaskScheduler: def _run_fetch_trends(self): """定时刷新热点趋势(百度/微博/知乎实时热搜 + LLM补充)""" + started = datetime.now(timezone.utc) + _log_task("scheduled_fetch_trends", "running", started_at=started) try: logger.info("[Scheduled] Fetching hot trends...") import subprocess @@ -100,14 +152,26 @@ class TaskScheduler: for line in result.stdout.strip().split("\n"): if line.strip(): logger.info("[Trends] %s", line.strip()) - logger.info("[Scheduled] Trends refreshed successfully") + _log_task("scheduled_fetch_trends", "success", + message="趋势刷新成功", + result_data={"output_lines": len(result.stdout.splitlines())}, + started_at=started, finished_at=datetime.now(timezone.utc)) else: - logger.warning("[Scheduled] Trends refresh failed: %s", result.stderr[-500:]) + _log_task("scheduled_fetch_trends", "failed", + message=f"返回码 {result.returncode}", + error_trace=result.stderr[-500:], + started_at=started, finished_at=datetime.now(timezone.utc)) except Exception as e: + _log_task("scheduled_fetch_trends", "failed", + message=str(e), + error_trace=traceback.format_exc(), + started_at=started, finished_at=datetime.now(timezone.utc)) logger.exception("[Scheduled] Trends refresh error: %s", e) def _run_refresh_search_cache(self): """定时刷新搜索缓存(通过 opencode webfetch)""" + started = datetime.now(timezone.utc) + _log_task("scheduled_refresh_search_cache", "running", started_at=started) try: logger.info("[Scheduled] Refreshing search cache via opencode...") import subprocess @@ -122,48 +186,95 @@ class TaskScheduler: if line.strip(): logger.warning("[SearchCache] %s", line.strip()) if result.returncode == 0: + _log_task("scheduled_refresh_search_cache", "success", + message="搜索缓存刷新成功", + result_data={"output_lines": len(result.stdout.splitlines())}, + started_at=started, finished_at=datetime.now(timezone.utc)) logger.info("[Scheduled] Search cache refreshed") else: + _log_task("scheduled_refresh_search_cache", "failed", + message="部分失败", + error_trace=result.stderr[-500:], + started_at=started, finished_at=datetime.now(timezone.utc)) logger.warning("[Scheduled] Search cache refresh may have partial failures") except subprocess.TimeoutExpired: + _log_task("scheduled_refresh_search_cache", "failed", + message="超时", + started_at=started, finished_at=datetime.now(timezone.utc)) logger.warning("[Scheduled] Search cache refresh timed out") except Exception as e: + _log_task("scheduled_refresh_search_cache", "failed", + message=str(e), + error_trace=traceback.format_exc(), + started_at=started, finished_at=datetime.now(timezone.utc)) logger.exception("[Scheduled] Search cache refresh error: %s", e) def _run_generate(self): + started = datetime.now(timezone.utc) + _log_task("scheduled_generate", "running", started_at=started) try: logger.info("[Scheduled] Starting content generation...") result = run_creator_blocking() logger.info("[Scheduled] Generation completed: %s", result) created_id = result.get("topic_id") if isinstance(result, dict) else None + review_result = None if created_id: logger.info("[Scheduled] Running compliance review on %s...", created_id) review_result = run_optimizer_blocking([created_id]) if review_result.get("ok"): logger.info("[Scheduled] Review completed for %s", created_id) - else: - logger.warning("[Scheduled] Review failed: %s", review_result.get("error")) + _log_task("scheduled_generate", "success", + message=f"创作完成" + (f", 选题 {created_id}" if created_id else ""), + result_data={"topic_id": created_id, "review_ok": review_result.get("ok") if review_result else None}, + started_at=started, finished_at=datetime.now(timezone.utc)) except Exception as e: + _log_task("scheduled_generate", "failed", + message=str(e), + error_trace=traceback.format_exc(), + started_at=started, finished_at=datetime.now(timezone.utc)) logger.exception("[Scheduled] Generation pipeline failed: %s", e) def _run_optimize(self): + started = datetime.now(timezone.utc) + _log_task("scheduled_optimize", "running", started_at=started) try: logger.info("[Scheduled] Starting compliance review...") result = run_optimizer_blocking() + _log_task("scheduled_optimize", "success", + message="合规审查完成", + result_data={"processed": result.get("processed", 0), "passed": result.get("passed", 0)}, + started_at=started, finished_at=datetime.now(timezone.utc)) logger.info("[Scheduled] Review completed: %s", result) except Exception as e: + _log_task("scheduled_optimize", "failed", + message=str(e), + error_trace=traceback.format_exc(), + started_at=started, finished_at=datetime.now(timezone.utc)) logger.exception("[Scheduled] Review failed: %s", e) def _run_collect(self): + started = datetime.now(timezone.utc) + _log_task("scheduled_collect", "running", started_at=started) try: logger.info("[Scheduled] Starting topic collection...") result = run_collector_blocking() + topics_count = result.get("topics_count", 0) + _log_task("scheduled_collect", "success", + message=f"采集完成,找到 {topics_count} 个选题", + result_data={"topics_count": topics_count, "output": str(result.get("output", ""))[:200]}, + started_at=started, finished_at=datetime.now(timezone.utc)) logger.info("[Scheduled] Collection completed: %s", result.get("output", "")[-200:]) except Exception as e: + _log_task("scheduled_collect", "failed", + message=str(e), + error_trace=traceback.format_exc(), + started_at=started, finished_at=datetime.now(timezone.utc)) logger.exception("[Scheduled] Collection failed: %s", e) def _run_optimize_sources(self): """AI自动优化采集类别与信息源:对比市场热点和当前配置,给出调整建议""" + started = datetime.now(timezone.utc) + _log_task("scheduled_optimize_sources", "running", started_at=started) try: logger.info("[Scheduled] Starting source optimization with AI...") from .nvidia_client import call_llm @@ -183,32 +294,16 @@ class TaskScheduler: cat_names = [c.name for c in cats] src_summary = "\n".join(f"- [{s.source_type}] {s.name}: {s.query or s.url or ''}" for s in sources) - prompt = f"""你是一个内容策略分析师。分析当前中文互联网可持续生活领域的真实热点,与以下配置进行对比。 + prompt = get_prompt("sources_optimization", + n=len(cat_names), + cat_names="\n".join(f"- {n}" for n in cat_names), + n2=len(sources), + src_summary=src_summary, + year=datetime.now().year, + ) -当前配置的类别({len(cat_names)}个): -{chr(10).join(f'- {n}' for n in cat_names)} - -当前配置的信息源({len(sources)}个): -{src_summary} - -请完成以下任务: -1. 评估每个类别是否仍符合2026年中国市场真实热点(基于你的知识) -2. 评估每个信息源是否可能在中国正常访问 -3. 建议新增或删除的类别(最多2条) -4. 建议新增的信息源搜索词(最多3条,包含具体搜索词) - -输出 JSON 格式: -{{ - "category_assessment": [{{"name": "类别名", "status": "保留/淘汰/合并", "reason": "原因"}}], - "source_assessment": [{{"name": "源名", "status": "保留/淘汰/替换", "reason": "原因"}}], - "suggested_new_categories": [{{"name": "类别名", "search_query": "搜索词", "reason": "推荐原因"}}], - "suggested_new_sources": [{{"name": "源名", "type": "web_search", "query": "搜索词", "focus": "聚焦领域"}}], - "summary": "一句话总结本次优化建议" -}} - -只输出JSON,不要其他文字。""" - - resp = call_llm(prompt, temperature=0.5, max_tokens=2000) + params = get_prompt_params("sources_optimization") + resp = call_llm(prompt, temperature=params.get("temperature", 0.5), max_tokens=params.get("max_tokens", 3000)) if resp.startswith("```"): resp = resp.split("\n", 1)[1].rsplit("\n", 1)[0] result = json.loads(resp) @@ -222,12 +317,25 @@ class TaskScheduler: db.add(SystemConfig(key="collector_ai_advice", value=json.dumps(result, ensure_ascii=False), description="AI每日采集优化建议")) db.commit() logger.info("[Scheduled] Source AI optimization completed: %s", result.get("summary", "")) + _log_task("scheduled_optimize_sources", "success", + message=result.get("summary", "优化完成"), + result_data={"categories_assessed": len(result.get("category_assessment", [])), + "sources_assessed": len(result.get("source_assessment", [])), + "suggested_cats": len(result.get("suggested_new_categories", [])), + "suggested_srcs": len(result.get("suggested_new_sources", []))}, + started_at=started, finished_at=datetime.now(timezone.utc)) db.close() except Exception as e: + _log_task("scheduled_optimize_sources", "failed", + message=str(e), + error_trace=traceback.format_exc(), + started_at=started, finished_at=datetime.now(timezone.utc)) logger.exception("[Scheduled] Source AI optimization failed: %s", e) def _run_metrics_sync(self): """定时从各平台公开API获取发布文章的效果数据(当前仅支持知乎)""" + started = datetime.now(timezone.utc) + _log_task("scheduled_metrics_sync", "running", started_at=started) try: logger.info("[Scheduled] Starting metrics sync (zhihu auto-fetch)...") from ..database import SessionLocal @@ -286,6 +394,10 @@ class TaskScheduler: if count: db.commit() logger.info("[Scheduled] Metrics sync completed: synced %d zhihu articles", count) + _log_task("scheduled_metrics_sync", "success", + message=f"同步完成,{count} 篇知乎文章", + result_data={"articles_synced": count}, + started_at=started, finished_at=datetime.now(timezone.utc)) # 生成指标反馈:按 field 聚合表现,写入 metrics_feedback.json 供 collector 读取 try: import json as json_mod @@ -317,9 +429,16 @@ class TaskScheduler: except Exception as e_fb: logger.warning("[Scheduled] Metrics feedback generation failed: %s", e_fb) else: + _log_task("scheduled_metrics_sync", "success", + message="无已发布的知乎文章", + started_at=started, finished_at=datetime.now(timezone.utc)) logger.info("[Scheduled] Metrics sync: no zhihu articles to sync") db.close() except Exception as e: + _log_task("scheduled_metrics_sync", "failed", + message=str(e), + error_trace=traceback.format_exc(), + started_at=started, finished_at=datetime.now(timezone.utc)) logger.exception("[Scheduled] Metrics sync failed: %s", e) def get_jobs(self): diff --git a/platform/backend/app/database.py b/platform/backend/app/database.py index 0b037f0..a9a1f48 100644 --- a/platform/backend/app/database.py +++ b/platform/backend/app/database.py @@ -61,6 +61,49 @@ def init_db(): ("platform_configs", "min_words", "INTEGER DEFAULT 0"), ("platform_configs", "max_words", "INTEGER DEFAULT 0"), ("platform_configs", "website_url", "VARCHAR"), + ("task_logs", "module_id", "VARCHAR"), + ("task_logs", "error_trace", "TEXT"), + ("task_logs", "triggered_by", "VARCHAR DEFAULT 'scheduler'"), + ("task_logs", "result_data", "JSON DEFAULT '{}'::json"), + ("task_logs", "next_run_time", "TIMESTAMP"), + ("task_configs", "module_id", "VARCHAR UNIQUE"), + ("task_configs", "enabled", "BOOLEAN DEFAULT TRUE"), + ("task_configs", "params", "JSON DEFAULT '{}'::json"), + ("task_configs", "schedule", "VARCHAR"), + ("task_configs", "last_modified_by", "VARCHAR"), + ("prompt_configs", "key", "VARCHAR UNIQUE"), + ("prompt_configs", "module_id", "VARCHAR"), + ("prompt_configs", "category", "VARCHAR DEFAULT 'prompt'"), + ("prompt_configs", "version", "VARCHAR DEFAULT 'v1'"), + ("prompt_configs", "content", "TEXT"), + ("prompt_configs", "variables", "JSON DEFAULT '[]'::json"), + ("prompt_configs", "description", "VARCHAR"), + ("prompt_configs", "enabled", "BOOLEAN DEFAULT TRUE"), + ("prompt_configs", "temperature", "FLOAT"), + ("prompt_configs", "max_tokens", "INTEGER"), + ("prompt_configs", "created_by", "VARCHAR"), + ("keyword_domain_map", "id", "INTEGER PRIMARY KEY"), + ("keyword_domain_map", "pattern", "VARCHAR"), + ("keyword_domain_map", "domain", "VARCHAR"), + ("keyword_domain_map", "sort_order", "INTEGER DEFAULT 0"), + ("keyword_domain_map", "is_active", "BOOLEAN DEFAULT TRUE"), + ("sensitive_words", "id", "INTEGER PRIMARY KEY"), + ("sensitive_words", "word", "VARCHAR"), + ("sensitive_words", "category", "VARCHAR DEFAULT 'general'"), + ("sensitive_words", "is_active", "BOOLEAN DEFAULT TRUE"), + ("sensitive_words", "added_by", "VARCHAR"), + ("content_clean_rules", "id", "INTEGER PRIMARY KEY"), + ("content_clean_rules", "rule_type", "VARCHAR"), + ("content_clean_rules", "pattern", "TEXT"), + ("content_clean_rules", "description", "VARCHAR"), + ("content_clean_rules", "is_active", "BOOLEAN DEFAULT TRUE"), + ("content_clean_rules", "sort_order", "INTEGER DEFAULT 0"), + ("collector_categories", "pain_template", "TEXT"), + ("trend_field_mappings", "id", "INTEGER PRIMARY KEY"), + ("trend_field_mappings", "trend_keyword", "VARCHAR"), + ("trend_field_mappings", "field_name", "VARCHAR"), + ("trend_field_mappings", "sort_order", "INTEGER DEFAULT 0"), + ("trend_field_mappings", "is_active", "BOOLEAN DEFAULT TRUE"), ]: try: conn.execute(text(f"ALTER TABLE {table} ADD COLUMN IF NOT EXISTS {col} {typ}")) diff --git a/platform/backend/app/main.py b/platform/backend/app/main.py index 1f2d039..b85ec15 100644 --- a/platform/backend/app/main.py +++ b/platform/backend/app/main.py @@ -9,7 +9,7 @@ from pathlib import Path from .database import engine, get_db, init_db from .models import Base -from .api import topics, system, articles, publishing, auth, admin, audit, optimizer_logs, cases, task_logs, llm_configs, system_configs, topic_config, calendar, metrics, assets, tasks, platform_config, collector_mgmt, assistant +from .api import topics, system, articles, publishing, auth, admin, audit, optimizer_logs, cases, task_logs, task_configs, prompt_configs, llm_configs, system_configs, topic_config, calendar, metrics, assets, tasks, platform_config, collector_mgmt, assistant, config_items from .initial_data import import_initial_data from .core.scheduler import scheduler @@ -83,6 +83,8 @@ app.include_router(admin.router) app.include_router(audit.router) app.include_router(optimizer_logs.router) app.include_router(cases.router) +app.include_router(task_configs.router) +app.include_router(prompt_configs.router) app.include_router(task_logs.router) app.include_router(llm_configs.router) app.include_router(system_configs.router) @@ -94,6 +96,7 @@ app.include_router(tasks.router) app.include_router(platform_config.router) app.include_router(collector_mgmt.router) app.include_router(assistant.router) +app.include_router(config_items.router) # 挂载自动生成的图片(必须先于前端根挂载) PROJECT_ROOT_DIR = Path(__file__).parent.parent.parent.parent diff --git a/platform/backend/app/models.py b/platform/backend/app/models.py index 5cf5052..1efb5dc 100644 --- a/platform/backend/app/models.py +++ b/platform/backend/app/models.py @@ -471,24 +471,59 @@ class TaskLog(Base): __tablename__ = "task_logs" id = Column(Integer, primary_key=True, index=True, autoincrement=True) + module_id = Column(String, nullable=False, index=True) # scheduled_collect / scheduled_generate 等 task_name = Column(String, nullable=False) - topic_id = Column(String, nullable=True) - status = Column(String, nullable=False) + topic_id = Column(String, nullable=True, index=True) + status = Column(String, nullable=False) # pending / running / success / failed / cancelled message = Column(Text, nullable=True) + error_trace = Column(Text, nullable=True) + triggered_by = Column(String, default="scheduler") # scheduler / manual / api + result_data = Column(JSON, default=dict) # 产出摘要:{topics_found, articles_created, issues_fixed, ...} started_at = Column(DateTime(timezone=True), server_default=func.now()) finished_at = Column(DateTime(timezone=True), nullable=True) - duration = Column(Integer, nullable=True) + duration = Column(Integer, nullable=True) # seconds + next_run_time = Column(DateTime(timezone=True), nullable=True) def to_dict(self): return { "id": self.id, + "module_id": self.module_id, "task_name": self.task_name, "topic_id": self.topic_id, "status": self.status, "message": self.message, + "error_trace": self.error_trace, + "triggered_by": self.triggered_by, + "result_data": self.result_data or {}, "started_at": self.started_at.isoformat() if self.started_at else None, "finished_at": self.finished_at.isoformat() if self.finished_at else None, "duration": self.duration, + "next_run_time": self.next_run_time.isoformat() if self.next_run_time else None, + } + + +class TaskConfig(Base): + __tablename__ = "task_configs" + + id = Column(Integer, primary_key=True, index=True, autoincrement=True) + module_id = Column(String, unique=True, nullable=False) + enabled = Column(Boolean, default=True) + params = Column(JSON, default=dict) # 各任务自定义参数,JSON 格式 + schedule = Column(String, nullable=True) # cron 表达式,覆盖默认 + last_modified_by = Column(String, nullable=True) + created_at = Column(DateTime(timezone=True), server_default=func.now()) + updated_at = Column(DateTime(timezone=True), server_default=func.now(), onupdate=func.now()) + + def to_dict(self): + return { + "id": self.id, + "module_id": self.module_id, + "enabled": self.enabled, + "params": self.params or {}, + "schedule": self.schedule, + "last_modified_by": self.last_modified_by, + "created_at": self.created_at.isoformat() if self.created_at else None, + "updated_at": self.updated_at.isoformat() if self.updated_at else None, } @@ -527,6 +562,101 @@ class LLMConfig(Base): } +class PromptConfig(Base): + __tablename__ = "prompt_configs" + + id = Column(Integer, primary_key=True, index=True, autoincrement=True) + key = Column(String, unique=True, nullable=False, index=True) + module_id = Column(String, nullable=True, index=True) + category = Column(String, default="prompt") # prompt / rule / template + version = Column(String, default="v1") + content = Column(Text, nullable=False) + variables = Column(JSON, default=[]) # [{name, description, default_value}] + description = Column(String, nullable=True) + enabled = Column(Boolean, default=True) + temperature = Column(Float, nullable=True) + max_tokens = Column(Integer, nullable=True) + created_by = Column(String, nullable=True) + created_at = Column(DateTime(timezone=True), server_default=func.now()) + updated_at = Column(DateTime(timezone=True), server_default=func.now(), onupdate=func.now()) + + def to_dict(self): + return { + "id": self.id, + "key": self.key, + "module_id": self.module_id, + "category": self.category, + "version": self.version, + "content": self.content, + "variables": self.variables or [], + "description": self.description, + "enabled": self.enabled, + "temperature": self.temperature, + "max_tokens": self.max_tokens, + "created_by": self.created_by, + "created_at": self.created_at.isoformat() if self.created_at else None, + "updated_at": self.updated_at.isoformat() if self.updated_at else None, + } + + +class KeywordDomainMap(Base): + __tablename__ = "keyword_domain_map" + + id = Column(Integer, primary_key=True, index=True, autoincrement=True) + pattern = Column(String, nullable=False) + domain = Column(String, nullable=False) + sort_order = Column(Integer, default=0) + is_active = Column(Boolean, default=True) + created_at = Column(DateTime(timezone=True), server_default=func.now()) + updated_at = Column(DateTime(timezone=True), onupdate=func.now()) + + def to_dict(self): + return { + "id": self.id, "pattern": self.pattern, "domain": self.domain, + "sort_order": self.sort_order, "is_active": self.is_active, + "created_at": self.created_at.isoformat() if self.created_at else None, + "updated_at": self.updated_at.isoformat() if self.updated_at else None, + } + + +class SensitiveWord(Base): + __tablename__ = "sensitive_words" + + id = Column(Integer, primary_key=True, index=True, autoincrement=True) + word = Column(String, nullable=False) + category = Column(String, default="general") + is_active = Column(Boolean, default=True) + added_by = Column(String, nullable=True) + created_at = Column(DateTime(timezone=True), server_default=func.now()) + + def to_dict(self): + return { + "id": self.id, "word": self.word, "category": self.category, + "is_active": self.is_active, "added_by": self.added_by, + "created_at": self.created_at.isoformat() if self.created_at else None, + } + + +class ContentCleanRule(Base): + __tablename__ = "content_clean_rules" + + id = Column(Integer, primary_key=True, index=True, autoincrement=True) + rule_type = Column(String, nullable=False) # thinking / preface / verbosity / html_thinking + pattern = Column(Text, nullable=False) + description = Column(String, nullable=True) + is_active = Column(Boolean, default=True) + sort_order = Column(Integer, default=0) + created_at = Column(DateTime(timezone=True), server_default=func.now()) + + def to_dict(self): + return { + "id": self.id, "rule_type": self.rule_type, "pattern": self.pattern, + "description": self.description, "is_active": self.is_active, + "sort_order": self.sort_order, + "created_at": self.created_at.isoformat() if self.created_at else None, + } + + class SystemConfig(Base): __tablename__ = "system_configs" @@ -557,6 +687,7 @@ class CollectorCategory(Base): name = Column(String, unique=True, nullable=False) description = Column(Text, nullable=True) search_query = Column(String, nullable=True) + pain_template = Column(Text, nullable=True) sort_order = Column(Integer, default=0) is_active = Column(Boolean, default=True) created_at = Column(DateTime(timezone=True), server_default=func.now()) @@ -570,6 +701,7 @@ class CollectorCategory(Base): "name": self.name, "description": self.description, "search_query": self.search_query, + "pain_template": self.pain_template, "sort_order": self.sort_order, "is_active": self.is_active, "created_at": self.created_at.isoformat() if self.created_at else None, @@ -577,6 +709,24 @@ class CollectorCategory(Base): } +class TrendFieldMapping(Base): + __tablename__ = "trend_field_mappings" + + id = Column(Integer, primary_key=True, index=True, autoincrement=True) + trend_keyword = Column(String, nullable=False) + field_name = Column(String, nullable=False) + sort_order = Column(Integer, default=0) + is_active = Column(Boolean, default=True) + created_at = Column(DateTime(timezone=True), server_default=func.now()) + + def to_dict(self): + return { + "id": self.id, "trend_keyword": self.trend_keyword, "field_name": self.field_name, + "sort_order": self.sort_order, "is_active": self.is_active, + "created_at": self.created_at.isoformat() if self.created_at else None, + } + + class CollectorSource(Base): """采集信息源(可在运营管理中动态编辑)""" __tablename__ = "collector_sources" diff --git a/platform/backend/app/schemas.py b/platform/backend/app/schemas.py index d302fca..bbca844 100644 --- a/platform/backend/app/schemas.py +++ b/platform/backend/app/schemas.py @@ -458,13 +458,18 @@ class CaseResponse(CaseBase): class TaskLogBase(BaseModel): + module_id: str task_name: str topic_id: Optional[str] = None status: str message: Optional[str] = None + error_trace: Optional[str] = None + triggered_by: str = "scheduler" + result_data: Dict[str, Any] = {} started_at: Optional[datetime] = None finished_at: Optional[datetime] = None duration: Optional[int] = None + next_run_time: Optional[datetime] = None class TaskLogResponse(TaskLogBase): @@ -473,6 +478,29 @@ class TaskLogResponse(TaskLogBase): model_config = ConfigDict(from_attributes=True) +class TaskConfigBase(BaseModel): + module_id: str + enabled: bool = True + params: Dict[str, Any] = {} + schedule: Optional[str] = None + + +class TaskConfigUpdate(BaseModel): + enabled: Optional[bool] = None + params: Optional[Dict[str, Any]] = None + schedule: Optional[str] = None + last_modified_by: Optional[str] = None + + +class TaskConfigResponse(TaskConfigBase): + id: int + last_modified_by: Optional[str] = None + created_at: Optional[datetime] = None + updated_at: Optional[datetime] = None + + model_config = ConfigDict(from_attributes=True) + + class LLMConfigBase(BaseModel): name: str system_prompt: Optional[str] = None diff --git a/platform/frontend/admin.html b/platform/frontend/admin.html index 9f58dce..c5491a1 100644 --- a/platform/frontend/admin.html +++ b/platform/frontend/admin.html @@ -22,65 +22,70 @@
{{ drawerData.description }}