diff --git a/automation/data/sustainability_cases.json b/automation/data/sustainability_cases.json index c664632..bfadd0d 100644 --- a/automation/data/sustainability_cases.json +++ b/automation/data/sustainability_cases.json @@ -1,4 +1,106 @@ [ + { + "id": "SUS-0A6B894E", + "country": "US", + "category": "循环消费", + "title": "8点1氪丨英伟达Q1净利润583亿美元;谷歌CEO:Gemini月活跃用户达9亿;寿司郎回应“盘子抽检10抽10脏”", + "core_idea": "

今日热点导览

\n

俄罗斯将延长对中国公民的免签制度

\n

张雪首度谈及签约德比斯原因

\n

微信再度开放520大额红包,限时1天 

\n

罗欣药业回应阿奇霉素劣药事件

\n

段永平最新持仓公布

\n

TOP3大新闻

\n

英伟达:第一财季净利润583亿美元,同比增长211%", + "data_facts": "数据点: 211, 85, 211", + "global_advantage": "需进一步分析全球优势", + "china_pain_point": "以旧换新流程繁琐、二手商品信任缺失、租赁市场不规范", + "localization_suggestion": "需基于中国现实调整实施", + "mvp_action": "建议先小规模试点验证", + "source_url": "https://36kr.com/p/3818280989443202?f=rss", + "credibility_rating": "⭐", + "china_applicability": "⭐⭐", + "collection_date": "2026-05-21", + "status": "待验证" + }, + { + "id": "SUS-CBB0B317", + "country": "US", + "category": "循环消费", + "title": "氪星晚报|2026福布斯中国人工智能科技企业TOP 50发布,中关村科金入选;三星劳资谈判破裂,威刚董事长称DRAM、NAND恐再掀涨价潮;新交所与中国银行签署新版战略合作备忘录", + "core_idea": "

大公司:

\n

新加坡加码人工智能布局,英伟达将落地本地研发中心

\n

全球AI芯片巨头英伟达", + "data_facts": "数据点: 100, 15, 30", + "global_advantage": "需进一步分析全球优势", + "china_pain_point": "以旧换新流程繁琐、二手商品信任缺失、租赁市场不规范", + "localization_suggestion": "需基于中国现实调整实施", + "mvp_action": "建议先小规模试点验证", + "source_url": "https://36kr.com/p/3816942423000196?f=rss", + "credibility_rating": "⭐", + "china_applicability": "⭐⭐", + "collection_date": "2026-05-21", + "status": "待验证" + }, + { + "id": "SUS-AEBB9FA8", + "country": "Global", + "category": "循环消费", + "title": "长鑫科技更新招股书,碧桂园错失300亿", + "core_idea": "

地产行业内,曾有望跑出两家明星创投企业,一家是世纪金源,另一家是碧桂园。前者现在做得顺风顺水,收获一个又一个IPO,而碧桂园因为地产业务暴雷,需要大量资金保交楼,被迫中断了在科技领域的投资“大爆发”。

\n

01

\n

5月17日,计划登陆科创板的长鑫科技集团股份有限公司(以下简称:长鑫科技)更新了招股书。新招股书的财务数据显示,", + "data_facts": "数据点: 719.13, 1688, 90", + "global_advantage": "需进一步分析全球优势", + "china_pain_point": "以旧换新流程繁琐、二手商品信任缺失、租赁市场不规范", + "localization_suggestion": "需基于中国现实调整实施", + "mvp_action": "建议先小规模试点验证", + "source_url": "https://36kr.com/p/3817039204139908?f=rss", + "credibility_rating": "⭐", + "china_applicability": "⭐⭐", + "collection_date": "2026-05-21", + "status": "待验证" + }, + { + "id": "SUS-15BBBA59", + "country": "Global", + "category": "循环消费", + "title": "要做数字劳动力生产工厂,「未来式智能」完成Pre-A轮融资", + "core_idea": "

文|王欣逸

\n

编辑|邓咏仪

\n

36氪获悉,未来式智能(AutoAgents.ai)近日完成Pre-A轮融资,新进投资方包括凡创资本、中关村资本、探元资本,老股东东证创新、麟阁创投跟投,本轮融资主要用于算力投入、团队扩张以及新产品的生态建设运营。

\n

未来式智能成立于2023年6月", + "data_facts": "数据点: 100, 72, 91", + "global_advantage": "需进一步分析全球优势", + "china_pain_point": "以旧换新流程繁琐、二手商品信任缺失、租赁市场不规范", + "localization_suggestion": "需基于中国现实调整实施", + "mvp_action": "建议先小规模试点验证", + "source_url": "https://36kr.com/p/3816909945029760?f=rss", + "credibility_rating": "⭐", + "china_applicability": "⭐⭐", + "collection_date": "2026-05-21", + "status": "待验证" + }, + { + "id": "SUS-265BE035", + "country": "US", + "category": "循环消费", + "title": "早期项目 | Chance AI获美图等数百万美元投资,用户数已达20万", + "core_idea": "

作者丨欧雪

\n

编辑丨袁斯来

\n

硬氪获悉,Visual Agent(视觉智能体)创业公司Chance AI宣布完成数百万美元天使轮融资。本轮由美图领投,NYX Ventures、阿里系投资机构等跟投。融资将用于模型能力迭代、北美学生群体增长、以及社区与商业化探索。

\n

Chance AI成立于2025年,核心产品是全球第一款以摄像头为主要交互入口的", + "data_facts": "数据点: 86.07, 85.4, 40", + "global_advantage": "需进一步分析全球优势", + "china_pain_point": "以旧换新流程繁琐、二手商品信任缺失、租赁市场不规范", + "localization_suggestion": "需基于中国现实调整实施", + "mvp_action": "建议先小规模试点验证", + "source_url": "https://36kr.com/p/3816947890471812?f=rss", + "credibility_rating": "⭐", + "china_applicability": "⭐⭐", + "collection_date": "2026-05-21", + "status": "待验证" + }, + { + "id": "CHN-001", + "country": "China", + "category": "循环消费", + "title": "中国二手交易平台崛起:闲鱼转转让闲置物品年交易额超5000亿", + "core_idea": "通过二手交易平台,用户可以将闲置物品变现,降低消费成本", + "data_facts": "闲鱼年交易额超5000亿,用户数超3亿,每天上架商品超200万件", + "global_advantage": "中国移动互联网普及率高,二手交易习惯逐渐养成", + "china_pain_point": "信任机制不完善,假货和退换货纠纷多", + "localization_suggestion": "选择信誉高的卖家,优先购买有质检服务的商品", + "mvp_action": "整理家中闲置物品,本月在二手平台卖出3件", + "source_url": "https://www.goofish.com/", + "credibility_rating": "⭐⭐⭐⭐", + "china_applicability": "⭐⭐⭐⭐⭐", + "collection_date": "2026-05-20", + "status": "已验证" + }, { "id": "GLO-002", "country": "Sweden", @@ -66,22 +168,5 @@ "china_applicability": "⭐⭐⭐", "collection_date": "2026-04-19", "status": "已验证" - }, - { - "id": "CHN-001", - "country": "China", - "category": "循环消费", - "title": "中国二手交易平台崛起:闲鱼转转让闲置物品年交易额超5000亿", - "core_idea": "通过二手交易平台,用户可以将闲置物品变现,降低消费成本", - "data_facts": "闲鱼年交易额超5000亿,用户数超3亿,每天上架商品超200万件", - "global_advantage": "中国移动互联网普及率高,二手交易习惯逐渐养成", - "china_pain_point": "信任机制不完善,假货和退换货纠纷多", - "localization_suggestion": "选择信誉高的卖家,优先购买有质检服务的商品", - "mvp_action": "整理家中闲置物品,本月在二手平台卖出3件", - "source_url": "https://www.goofish.com/", - "credibility_rating": "⭐⭐⭐⭐", - "china_applicability": "⭐⭐⭐⭐⭐", - "collection_date": "2026-05-20", - "status": "已验证" } ] \ No newline at end of file diff --git a/platform/backend/app/api/system.py b/platform/backend/app/api/system.py index a89e57b..e0cbc03 100644 --- a/platform/backend/app/api/system.py +++ b/platform/backend/app/api/system.py @@ -1,18 +1,19 @@ import logging +import subprocess from fastapi import APIRouter, HTTPException, Depends, Body from sqlalchemy.orm import Session from sqlalchemy import func from datetime import datetime, date from pathlib import Path from typing import Dict, Any, List, Optional -from pathlib import Path import os import json from ..database import get_db from ..models import Topic, Article -from ..core.generator import run_creator -from ..core.optimizer import run_optimizer -from ..core.collector import run_collector +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 +import threading from ..core.sync import sync_all_topics from ..core.scheduler import scheduler from .auth import get_current_user, org_filter @@ -61,59 +62,51 @@ 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)): - logger.info(f"Received topic_id={topic_id}") + logger.info(f"Generation triggered by {current_user.username}, topic_id={topic_id}") try: result = run_creator(topic_id) - if not result["ok"]: - raise HTTPException(status_code=500, detail=result["error"]) - sync_all_topics() - return {"message": "Generation triggered", "result": result} + return {"message": "内容创作已后台启动", "pid": result.get("pid")} except Exception as e: raise HTTPException(status_code=500, detail=str(e)) +@router.get("/generate/status", dependencies=[Depends(get_current_user)]) +def generation_status(): + status = get_generator_status() + if status is None: + return {"status": "idle", "message": "当前无运行中的创作任务"} + return status + @router.post("/collect/run", dependencies=[Depends(get_current_user)]) def trigger_collection(db: Session = Depends(get_db), current_user=Depends(get_current_user)): logger.info(f"Manual collection triggered by {current_user.username}") try: result = run_collector() - return {"message": "内容采集已完成", "result": result} + return {"message": "内容采集已后台启动", "result": result} except Exception as e: raise HTTPException(status_code=500, detail=str(e)) +@router.get("/collect/status", dependencies=[Depends(get_current_user)]) +def collection_status(): + status = get_collector_status() + if status is None: + return {"status": "idle", "message": "当前无运行中的采集任务"} + 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)): try: result = run_optimizer(topic_ids) - if not result["ok"]: - raise HTTPException(status_code=500, detail=result["error"]) - report = result.get("report") - if report and report["summary"]["total_articles"] > 0: - s = report["summary"] - total = s["total_articles"] - avg = s["average_score"] - msg = f"审查完成: {total} 篇全部通过 ({avg:.0f}分)" - return {"message": msg, "summary": s} - else: - if topic_ids: - updated = 0 - for tid in topic_ids: - q = db.query(Topic).filter(Topic.id == tid) - of = org_filter(current_user, Topic) - if of is not True: - q = q.filter(of) - topic = q.first() - if topic and topic.status in ('review', '待审查'): - topic.status = 'ready' - if not topic.generated_at: - topic.generated_at = datetime.utcnow() - updated += 1 - db.commit() - if updated: - logger.info(f"Review: {updated} topics advanced to 'ready' (no release files)") - return {"message": "审查完成(未找到 release 文件,仅推进状态)", "stdout": result.get("stdout", "")} + return {"message": "合规审查已后台启动", "pid": result.get("pid")} except Exception as e: raise HTTPException(status_code=500, detail=str(e)) +@router.get("/review/status", dependencies=[Depends(get_current_user)]) +def review_status(): + status = get_optimizer_status() + if status is None: + return {"status": "idle", "message": "当前无运行中的审查任务"} + return status + @router.get("/logs/{log_date}", dependencies=[Depends(get_current_user)]) def get_logs(log_date: str, log_type: str = "creator"): log_file = LOGS_DIR / f"{log_type}_{log_date}.log" @@ -157,11 +150,14 @@ def run_sync(): def trigger_optimize_sources(): try: from ..core.scheduler import scheduler - scheduler._run_optimize_sources() - log_file = LOGS_DIR / f"optimizer_sources_{date.today().isoformat()}.log" - log_file.parent.mkdir(parents=True, exist_ok=True) - log_file.write_text(f"{datetime.now().isoformat()} - 信息源优化完成\n") - return {"message": "信息源优化已完成"} + def _bg(): + try: + scheduler._run_optimize_sources() + except Exception as e: + logger.exception("Background optimize sources failed: %s", e) + t = threading.Thread(target=_bg, daemon=True) + t.start() + return {"message": "信息源优化已后台启动"} except Exception as e: raise HTTPException(status_code=500, detail=str(e)) @@ -169,44 +165,44 @@ def trigger_optimize_sources(): def trigger_metrics_sync(): try: from ..core.scheduler import scheduler - scheduler._run_metrics_sync() - log_file = LOGS_DIR / f"metrics_sync_{date.today().isoformat()}.log" - log_file.parent.mkdir(parents=True, exist_ok=True) - log_file.write_text(f"{datetime.now().isoformat()} - 指标同步完成\n") - return {"message": "指标同步已完成"} + def _bg(): + try: + scheduler._run_metrics_sync() + except Exception as e: + logger.exception("Background metrics sync failed: %s", e) + t = threading.Thread(target=_bg, daemon=True) + t.start() + return {"message": "指标同步已后台启动"} except Exception as e: raise HTTPException(status_code=500, detail=str(e)) @router.post("/refresh-search-cache/run") def trigger_refresh_search_cache(): - """手动刷新搜索缓存""" try: - import subprocess, sys as sys_mod + import sys as sys_mod scripts_dir = Path(__file__).parent.parent.parent.parent / "scripts" - result = subprocess.run( + proc = subprocess.Popen( [sys_mod.executable, str(scripts_dir / "opencode_search.py"), "--refresh-cache"], - capture_output=True, text=True, timeout=600 + stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, + cwd=scripts_dir.parent.parent ) - if result.returncode != 0: - raise Exception(result.stderr[-500:]) - return {"message": "搜索缓存已刷新", "output": result.stdout.strip()} + logger.info("Search cache refresh started (pid=%s)", proc.pid) + return {"message": "搜索缓存刷新已后台启动", "pid": proc.pid} except Exception as e: raise HTTPException(status_code=500, detail=str(e)) @router.post("/trends/run") def trigger_trends_refresh(): - """手动刷新热点趋势数据""" try: - import subprocess, sys as sys_mod - from pathlib import Path + import sys as sys_mod scripts_dir = Path(__file__).parent.parent.parent.parent / "scripts" - result = subprocess.run( + proc = subprocess.Popen( [sys_mod.executable, str(scripts_dir / "trends.py")], - capture_output=True, text=True, timeout=120 + stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, + cwd=scripts_dir.parent.parent ) - if result.returncode != 0: - raise Exception(result.stderr[-500:]) - return {"message": "热点趋势已刷新", "output": result.stdout.strip()} + logger.info("Trends refresh started (pid=%s)", proc.pid) + return {"message": "热点趋势刷新已后台启动", "pid": proc.pid} except Exception as e: raise HTTPException(status_code=500, detail=str(e)) diff --git a/platform/backend/app/core/collector.py b/platform/backend/app/core/collector.py index 5d80062..c2e7bd6 100644 --- a/platform/backend/app/core/collector.py +++ b/platform/backend/app/core/collector.py @@ -1,36 +1,84 @@ """ 选题收集模块 调用 scripts/collector.py 脚本,从外部来源收集/生成新选题并写入 JSON +使用非阻塞 subprocess.Popen 避免阻塞 uvicorn worker """ import subprocess from pathlib import Path import logging import os +import time +from typing import Optional, Dict logger = logging.getLogger(__name__) -# 计算项目根目录 PROJECT_ROOT = Path(__file__).resolve().parents[4] if os.getenv('PROJECT_ROOT'): PROJECT_ROOT = Path(os.getenv('PROJECT_ROOT')) -def run_collector(): - """运行选题收集脚本""" +_running_processes: Dict[str, dict] = {} + +def _get_cmd(): script_path = PROJECT_ROOT / "scripts" / "collector.py" if not script_path.exists(): raise FileNotFoundError(f"Collector script not found: {script_path}") venv_python = PROJECT_ROOT / "platform" / "backend" / "venv" / "bin" / "python" if venv_python.exists(): - cmd = [str(venv_python), str(script_path)] - else: - cmd = ["python3", str(script_path)] + return [str(venv_python), str(script_path)] + return ["python3", str(script_path)] + +def run_collector(): + """运行选题收集脚本(非阻塞,后台运行)""" + cmd = _get_cmd() + proc = subprocess.Popen( + cmd, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + cwd=PROJECT_ROOT + ) + _running_processes["collector"] = { + "pid": proc.pid, + "started_at": time.time(), + "process": proc + } + logger.info("Collector started in background (pid=%s)", proc.pid) + return {"ok": True, "pid": proc.pid} + +def run_collector_blocking(timeout: int = 300): + """运行选题收集脚本(阻塞,带超时,给定时任务用)""" + cmd = _get_cmd() result = subprocess.run( cmd, capture_output=True, text=True, cwd=PROJECT_ROOT, - timeout=300 + timeout=timeout ) if result.returncode != 0: raise RuntimeError(f"Collector failed: {result.stderr}") - return {"ok": True, "output": result.stdout} \ No newline at end of file + return {"ok": True, "output": result.stdout} + +def get_collector_status() -> Optional[Dict]: + """获取当前采集任务状态""" + info = _running_processes.get("collector") + if not info: + return None + proc: subprocess.Popen = info["process"] + if proc.poll() is not None: + stdout, stderr = proc.communicate() + elapsed = time.time() - info["started_at"] + del _running_processes["collector"] + return { + "status": "completed" if proc.returncode == 0 else "failed", + "pid": info["pid"], + "elapsed": round(elapsed, 1), + "returncode": proc.returncode, + "stdout": stdout.strip()[-500:], + "stderr": stderr.strip()[-500:], + } + return { + "status": "running", + "pid": info["pid"], + "elapsed": round(time.time() - info["started_at"], 1), + } diff --git a/platform/backend/app/core/generator.py b/platform/backend/app/core/generator.py index fe3215a..5d9fb20 100644 --- a/platform/backend/app/core/generator.py +++ b/platform/backend/app/core/generator.py @@ -2,37 +2,40 @@ import subprocess from pathlib import Path import logging import os +import time +from typing import Optional, Dict logger = logging.getLogger(__name__) -# 计算项目根目录(从本文件位置上升4层) PROJECT_ROOT = Path(__file__).resolve().parents[4] -# 允许环境变量覆盖(适合容器部署) if os.getenv('PROJECT_ROOT'): PROJECT_ROOT = Path(os.getenv('PROJECT_ROOT')) -def run_creator(topic_id: str = None): - """运行内容创作脚本,返回简略结果 - - Args: - topic_id: 可选,指定要创作的选题ID。不指定则创作优先级最高的选题。 - """ - from datetime import datetime, timezone +_running_processes: Dict[str, dict] = {} + +def _get_cmd(topic_id: str = None): script_path = PROJECT_ROOT / "scripts" / "creator.py" + if not script_path.exists(): + raise FileNotFoundError(f"Creator script not found: {script_path}") venv_python = PROJECT_ROOT / "platform" / "backend" / "venv" / "bin" / "python" - if venv_python.exists(): - cmd = [str(venv_python), str(script_path)] - else: - cmd = ["python3", str(script_path)] + cmd = [str(venv_python), str(script_path)] if venv_python.exists() else ["python3", str(script_path)] if topic_id: cmd.extend(["--topic-id", topic_id]) - result = subprocess.run( - cmd, - cwd=str(PROJECT_ROOT), - capture_output=True, - text=True, - timeout=1800 - ) + return cmd + +def run_creator(topic_id: str = None): + """非阻塞:后台启动创作脚本,立即返回""" + cmd = _get_cmd(topic_id) + proc = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, cwd=str(PROJECT_ROOT)) + _running_processes["generator"] = {"pid": proc.pid, "started_at": time.time(), "process": proc, "topic_id": topic_id} + logger.info("Creator started in background (pid=%s, topic_id=%s)", proc.pid, topic_id) + return {"ok": True, "pid": proc.pid, "topic_id": topic_id} + +def run_creator_blocking(topic_id: str = None, timeout: int = 1800): + """阻塞版:带超时,给定时任务使用""" + from datetime import datetime, timezone + cmd = _get_cmd(topic_id) + result = subprocess.run(cmd, cwd=str(PROJECT_ROOT), capture_output=True, text=True, timeout=timeout) if result.returncode != 0: return {"ok": False, "error": result.stderr} @@ -55,9 +58,6 @@ def run_creator(topic_id: str = None): if topic.status in ('pending', '待处理'): topic.status = 'review' db.commit() - logger.info(f"Topic {topic_id} updated: generated_at set, status→{topic.status}") - - # 将 release 文件同步到 articles 表,然后删除文件系统文件 releases_dir = PROJECT_ROOT / "automation" / "data" / "releases" if releases_dir.exists(): for dd in sorted(releases_dir.iterdir(), reverse=True): @@ -74,17 +74,8 @@ def run_creator(topic_id: str = None): if existing: existing.html_content = html else: - db.add(Article( - id=article_id, - topic_id=topic_id, - platform=platform_dir, - file_path=f"db:{article_id}", - html_content=html, - status="draft", - )) + db.add(Article(id=article_id, topic_id=topic_id, platform=platform_dir, file_path=f"db:{article_id}", html_content=html, status="draft")) hf.unlink() - logger.info(f"Synced {hf.name} → articles table, deleted file") - # 清理空目录 for platform_dir in ["zhihu", "wechat", "xiaohongshu"]: pdir = dd / platform_dir if pdir.exists() and not any(pdir.iterdir()): @@ -95,8 +86,16 @@ def run_creator(topic_id: str = None): except Exception as e: logger.warning(f"DB sync after creation failed: {e}") - return { - "ok": True, - "topic_id": topic_id, - "stdout": result.stdout[-1000:] if len(result.stdout) > 1000 else result.stdout - } + return {"ok": True, "topic_id": topic_id, "stdout": result.stdout[-1000:] if len(result.stdout) > 1000 else result.stdout} + +def get_generator_status() -> Optional[Dict]: + info = _running_processes.get("generator") + if not info: + return None + proc = info["process"] + if proc.poll() is not None: + stdout, stderr = proc.communicate() + elapsed = time.time() - info["started_at"] + del _running_processes["generator"] + return {"status": "completed" if proc.returncode == 0 else "failed", "pid": info["pid"], "elapsed": round(elapsed, 1), "returncode": proc.returncode, "topic_id": info["topic_id"]} + return {"status": "running", "pid": info["pid"], "elapsed": round(time.time() - info["started_at"], 1)} diff --git a/platform/backend/app/core/optimizer.py b/platform/backend/app/core/optimizer.py index 863b046..17b26ea 100644 --- a/platform/backend/app/core/optimizer.py +++ b/platform/backend/app/core/optimizer.py @@ -3,44 +3,43 @@ from pathlib import Path import logging import os import json +import time from datetime import datetime -from typing import List +from typing import List, Optional, Dict logger = logging.getLogger(__name__) -# 计算项目根目录(从本文件位置上升4层) PROJECT_ROOT = Path(__file__).resolve().parents[4] if os.getenv('PROJECT_ROOT'): PROJECT_ROOT = Path(os.getenv('PROJECT_ROOT')) -def run_optimizer(topic_ids: List[str] = None): - """运行合规优化脚本,返回报告摘要 - - Args: - topic_ids: 可选,指定要优化的选题ID列表。不指定则优化所有 draft 文章。 - """ +_running_processes: Dict[str, dict] = {} + +def _get_cmd(topic_ids: List[str] = None): script_path = PROJECT_ROOT / "scripts" / "compliance_optimizer.py" + if not script_path.exists(): + raise FileNotFoundError(f"Optimizer script not found: {script_path}") venv_python = PROJECT_ROOT / "platform" / "backend" / "venv" / "bin" / "python" - if venv_python.exists(): - cmd = [str(venv_python), str(script_path)] - else: - cmd = ["python3", str(script_path)] + cmd = [str(venv_python), str(script_path)] if venv_python.exists() else ["python3", str(script_path)] if topic_ids: cmd.extend(["--topic-ids", ','.join(topic_ids)]) - logger.info(f"[DEBUG] Running optimizer with topic_ids={topic_ids}, cmd={' '.join(cmd)}") - - result = subprocess.run( - cmd, - cwd=str(PROJECT_ROOT), - capture_output=True, - text=True, - timeout=600 # 10分钟 - ) + return cmd + +def run_optimizer(topic_ids: List[str] = None): + """非阻塞:后台启动合规审查脚本,立即返回""" + cmd = _get_cmd(topic_ids) + proc = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, cwd=str(PROJECT_ROOT)) + _running_processes["optimizer"] = {"pid": proc.pid, "started_at": time.time(), "process": proc, "topic_ids": topic_ids} + logger.info("Optimizer started in background (pid=%s, topic_ids=%s)", proc.pid, topic_ids) + return {"ok": True, "pid": proc.pid, "topic_ids": topic_ids} + +def run_optimizer_blocking(topic_ids: List[str] = None, timeout: int = 600): + """阻塞版:带超时,给定时任务使用""" + cmd = _get_cmd(topic_ids) + result = subprocess.run(cmd, cwd=str(PROJECT_ROOT), capture_output=True, text=True, timeout=timeout) if result.returncode != 0: logger.error(f"Optimizer failed: {result.stderr}") return {"ok": False, "error": result.stderr} - - # 读取优化报告(优化脚本会在 today 的 drafts 目录生成报告) report_date = datetime.now().strftime("%Y-%m-%d") report_path = PROJECT_ROOT / "automation" / "data" / "drafts" / report_date / "optimization_report.json" if report_path.exists(): @@ -49,3 +48,15 @@ def run_optimizer(topic_ids: List[str] = None): else: logger.warning(f"Report not found: {report_path}") return {"ok": True, "report": None, "stdout": result.stdout} + +def get_optimizer_status() -> Optional[Dict]: + info = _running_processes.get("optimizer") + if not info: + return None + proc = info["process"] + if proc.poll() is not None: + stdout, stderr = proc.communicate() + elapsed = time.time() - info["started_at"] + del _running_processes["optimizer"] + return {"status": "completed" if proc.returncode == 0 else "failed", "pid": info["pid"], "elapsed": round(elapsed, 1), "returncode": proc.returncode, "stdout": stdout.strip()[-300:], "stderr": stderr.strip()[-300:]} + return {"status": "running", "pid": info["pid"], "elapsed": round(time.time() - info["started_at"], 1)} diff --git a/platform/backend/app/core/scheduler.py b/platform/backend/app/core/scheduler.py index a362306..d7e4440 100644 --- a/platform/backend/app/core/scheduler.py +++ b/platform/backend/app/core/scheduler.py @@ -9,9 +9,9 @@ from datetime import datetime from pathlib import Path from apscheduler.schedulers.background import BackgroundScheduler from apscheduler.triggers.cron import CronTrigger -from .generator import run_creator -from .optimizer import run_optimizer -from .collector import run_collector +from .generator import run_creator_blocking +from .optimizer import run_optimizer_blocking +from .collector import run_collector_blocking logger = logging.getLogger(__name__) @@ -132,12 +132,12 @@ class TaskScheduler: def _run_generate(self): try: logger.info("[Scheduled] Starting content generation...") - result = run_creator() + result = run_creator_blocking() logger.info("[Scheduled] Generation completed: %s", result) created_id = result.get("topic_id") if isinstance(result, dict) else None if created_id: logger.info("[Scheduled] Running compliance review on %s...", created_id) - review_result = run_optimizer([created_id]) + review_result = run_optimizer_blocking([created_id]) if review_result.get("ok"): logger.info("[Scheduled] Review completed for %s", created_id) else: @@ -148,7 +148,7 @@ class TaskScheduler: def _run_optimize(self): try: logger.info("[Scheduled] Starting compliance review...") - result = run_optimizer() + result = run_optimizer_blocking() logger.info("[Scheduled] Review completed: %s", result) except Exception as e: logger.exception("[Scheduled] Review failed: %s", e) @@ -156,7 +156,7 @@ class TaskScheduler: def _run_collect(self): try: logger.info("[Scheduled] Starting topic collection...") - result = run_collector() + result = run_collector_blocking() logger.info("[Scheduled] Collection completed: %s", result.get("output", "")[-200:]) except Exception as e: logger.exception("[Scheduled] Collection failed: %s", e) diff --git a/start-platform.sh b/start-platform.sh index a72c313..46fca70 100755 --- a/start-platform.sh +++ b/start-platform.sh @@ -46,4 +46,4 @@ echo "========================================" echo "" cd backend -exec $PYTHON_CMD -m uvicorn app.main:app --host 0.0.0.0 --port $PORT +exec $PYTHON_CMD -m uvicorn app.main:app --host 0.0.0.0 --port $PORT --workers 2