feat: pipeline 并行摘要 + 子进程 env 注入 + 循环导入修复
源码: - runtime/__init__.py: resume_jobs/resume_service 改为懒加载,打破循环导入 - freshrss_pipeline_jobs.py / resume_jobs.py: 子进程注入 .env 环境变量 - freshrss_pipeline.py: LLM 摘要串行改并行 (ThreadPoolExecutor, max_workers=4) 配置: - term_aliases: 19→149 条,大幅扩充别名映射 - term_stopwords: 19→132 条,增加过滤规则 - filter_context.personal.json: +7 个兴趣关键词
This commit is contained in:
@@ -33,8 +33,10 @@
|
||||
"Claude",
|
||||
"Claude Code",
|
||||
"CLAUDE.md",
|
||||
"CLI",
|
||||
"Context Engineering",
|
||||
"Cursor",
|
||||
"ChatGPT",
|
||||
"DeepSeek",
|
||||
"FastAPI",
|
||||
"Gin",
|
||||
@@ -46,7 +48,9 @@
|
||||
"Kafka",
|
||||
"Kubernetes",
|
||||
"LLM",
|
||||
"Loop Engineering",
|
||||
"MCP",
|
||||
"MoE",
|
||||
"MySQL",
|
||||
"MySQL复制延迟",
|
||||
"OpenAI",
|
||||
@@ -55,6 +59,7 @@
|
||||
"Prompt Engineering",
|
||||
"Python",
|
||||
"RAG",
|
||||
"ReAct",
|
||||
"ReActAgent",
|
||||
"Redis",
|
||||
"Skill",
|
||||
@@ -74,6 +79,8 @@
|
||||
"向量数据库",
|
||||
"多Agent协作",
|
||||
"大模型",
|
||||
"子Agent",
|
||||
"强化学习",
|
||||
"微服务",
|
||||
"渐进式披露",
|
||||
"知识库"
|
||||
|
||||
+142
-12
@@ -1,19 +1,149 @@
|
||||
{
|
||||
"AI助手": "AI Agent",
|
||||
"图文RAG": "RAG",
|
||||
"Prompt架构": "Prompt Engineering",
|
||||
"Agent Skill": "Agent Skills",
|
||||
"Binlog": "binlog",
|
||||
"Agent": "Agent",
|
||||
"Agent框架": "Agent",
|
||||
"Agent能力": "Agent Skills",
|
||||
"智能体": "AI Agent",
|
||||
"Agentic架构": "Agentic架构",
|
||||
"多Agent协作": "多Agent协作",
|
||||
"多智能体架构": "多Agent协作",
|
||||
"Multi-Agent": "多Agent",
|
||||
"Subagent": "子Agent",
|
||||
"Sub Agents验证": "子Agent",
|
||||
"子Agent": "子Agent",
|
||||
"子智能体": "子Agent",
|
||||
"Coding Agent": "AI Coding Agent",
|
||||
"Subagent": "SubAgent",
|
||||
"Subagents": "SubAgent",
|
||||
"vibe coding": "Vibe Coding",
|
||||
"AI编程": "AI Coding Agent",
|
||||
"AI辅助编程": "AI Coding Agent",
|
||||
"代码生成": "AI代码生成",
|
||||
"代码审查": "Code Review",
|
||||
"Prompt": "Prompt Engineering",
|
||||
"Prompt Caching": "提示缓存",
|
||||
"RAG": "RAG",
|
||||
"图文RAG": "RAG",
|
||||
"向量检索": "向量检索",
|
||||
"向量嵌入": "向量嵌入",
|
||||
"Multi-Token Prediction": "多Token预测",
|
||||
"Pair-In Pair-Out": "PIPO架构",
|
||||
"PIPO": "PIPO架构",
|
||||
"上下文管理": "上下文管理",
|
||||
"上下文卸载": "上下文卸载",
|
||||
"Self-GC": "上下文压缩",
|
||||
"记忆管理": "上下文管理",
|
||||
"会话管理": "上下文管理",
|
||||
"Harness Engineering": "Harness工程化",
|
||||
"Harness架构": "Harness工程化",
|
||||
"Harness": "Harness工程化",
|
||||
"Loop Engineering": "Loop Engineering",
|
||||
"推理加速": "推理加速",
|
||||
"推理深度": "推理深度",
|
||||
"长链路推理": "长链路推理",
|
||||
"RLVR": "RLVR",
|
||||
"GRPO": "GRPO",
|
||||
"强化学习": "强化学习",
|
||||
"Multi-Agent RL": "多Agent强化学习",
|
||||
"Viking AI搜索": "AI搜索",
|
||||
"Viking AI Search": "AI搜索",
|
||||
"智能搜索": "AI搜索",
|
||||
"SearchCLI": "CLI搜索",
|
||||
"视频生成": "AI视频生成",
|
||||
"视频生成模型": "AI视频生成",
|
||||
"LingBot-Video": "AI视频生成",
|
||||
"视觉自回归模型": "AI视频生成",
|
||||
"火山云数据库PostgreSQL Serverless版": "Serverless数据库",
|
||||
"PostgreSQL": "PostgreSQL",
|
||||
"MySQL": "MySQL",
|
||||
"OceanBase": "OceanBase",
|
||||
"StarRocks": "StarRocks",
|
||||
"Milvus": "Milvus",
|
||||
"Seal AI Zone": "AI安全",
|
||||
"NEX沙箱": "沙箱隔离",
|
||||
"MicroVM": "沙箱隔离",
|
||||
"安全左移": "安全左移",
|
||||
"安全中台": "AI安全",
|
||||
"成本降低": "成本优化",
|
||||
"成本杠杆": "成本优化",
|
||||
"Scale-to-Zero": "弹性伸缩",
|
||||
"Data as Git": "数据分支管理",
|
||||
"Schema Diff": "Schema对比",
|
||||
"Time Travel": "数据回溯",
|
||||
"多端架构": "多端架构",
|
||||
"契约化": "契约化架构",
|
||||
"大仓": "大仓工程化",
|
||||
"Vibe Coding": "Vibe Coding",
|
||||
"LLM Judge": "LLM评估",
|
||||
"SWE-Bench": "SWE-Bench",
|
||||
"SWE Bench Pro": "SWE-Bench",
|
||||
"SWE-Bench Pro": "SWE-Bench",
|
||||
"Verification Agent": "验证Agent",
|
||||
"CLI工具": "CLI",
|
||||
"CLI": "CLI",
|
||||
"漏桶算法": "限流架构",
|
||||
"固定窗口限流": "限流架构",
|
||||
"Suspend消费控制": "限流架构",
|
||||
"RocketMQ LiteTopic": "消息队列",
|
||||
"LLM Wiki": "LLM知识库",
|
||||
"知识工程": "知识工程",
|
||||
"语义资产": "语义资产管理",
|
||||
"知识图谱": "知识图谱",
|
||||
"知识库沉淀": "知识管理",
|
||||
"Skill": "Skill",
|
||||
"Skill Hub": "技能生态",
|
||||
"具身智能": "具身智能",
|
||||
"Open X-Embodiment": "具身智能",
|
||||
"YOLO Classifier": "目标检测",
|
||||
"MCP": "MCP",
|
||||
"MCP连接器": "MCP",
|
||||
"缓存击穿": "缓存优化",
|
||||
"GPU算力调度": "算力调度",
|
||||
"异构资源": "异构计算",
|
||||
"XPU": "异构计算",
|
||||
"弹性RDMA": "RDMA网络",
|
||||
"国内主流GPU": "国产芯片",
|
||||
"国产AI芯片": "国产芯片",
|
||||
"Paxos协议": "分布式一致性",
|
||||
"Token": "Token管理",
|
||||
"Token效率": "Token管理",
|
||||
"百万token上下文": "长上下文",
|
||||
"MoE": "MoE架构",
|
||||
"MoE架构": "MoE架构",
|
||||
"思维链": "思维链",
|
||||
"CoT Distillation": "思维链蒸馏",
|
||||
"自然语言驱动": "自然语言交互",
|
||||
"NL2SQL": "NL2SQL",
|
||||
"AI对齐": "AI对齐",
|
||||
"注意力机制": "注意力机制",
|
||||
"多模态": "多模态",
|
||||
"音视频工作台": "音视频处理",
|
||||
"AI助手": "AI Agent",
|
||||
"Agent架构": "AI Agent",
|
||||
"Agent专业化": "AI Agent",
|
||||
"Agent Teams": "多Agent协作",
|
||||
"Agentic Engineering": "AI Agent",
|
||||
"CLI工具": "CLI",
|
||||
"AI编程": "AI Coding Agent",
|
||||
"记忆管理": "上下文管理",
|
||||
"会话管理": "上下文管理"
|
||||
"AI智能体": "AI Agent",
|
||||
"LLM Agent": "AI Agent",
|
||||
"AI Harness": "Harness Engineering",
|
||||
"AI代码生成": "AI Coding Agent",
|
||||
"Memory管理": "上下文管理",
|
||||
"Agent Skill": "Agent Skills",
|
||||
"Binlog": "binlog",
|
||||
"vibe coding": "Vibe Coding",
|
||||
"Agent组织化协作平台": "Agent协作平台",
|
||||
"Anthropic": "Anthropic",
|
||||
"OpenClaw": "OpenClaw",
|
||||
"WorkBuddy": "WorkBuddy",
|
||||
"Claude": "Claude",
|
||||
"ChatGPT": "ChatGPT",
|
||||
"GPT": "GPT",
|
||||
"Opus": "Opus",
|
||||
"Sonnet": "Sonnet",
|
||||
"Grok": "Grok",
|
||||
"Qwen": "Qwen",
|
||||
"GLM": "GLM",
|
||||
"Claude Code": "Claude Code",
|
||||
"Cursor": "Cursor",
|
||||
"Codex": "Codex",
|
||||
"Pi": "Pi",
|
||||
"CoT": "CoT",
|
||||
"SVG": "SVG",
|
||||
"TTS": "TTS"
|
||||
}
|
||||
|
||||
@@ -536,6 +536,60 @@
|
||||
"value": "上下文管理",
|
||||
"suggestion_date": "2026-05-14",
|
||||
"based_on_days": 365
|
||||
},
|
||||
{
|
||||
"applied_at": "2026-07-15T02:17:50.155586Z",
|
||||
"action": "add_interest_keyword",
|
||||
"term": "强化学习",
|
||||
"reason": "top 0.8% by frequency,growth=100%,not yet covered by interest keywords or stopwords. Evidence: total_count=10, days_seen=10, recent_count=10.",
|
||||
"suggestions_path": "outputs/term_index/review/term-cleanup-suggestions-2026-07-15.json",
|
||||
"suggestion_date": "2026-07-15",
|
||||
"based_on_days": 365
|
||||
},
|
||||
{
|
||||
"applied_at": "2026-07-15T02:17:50.155586Z",
|
||||
"action": "add_interest_keyword",
|
||||
"term": "ReAct",
|
||||
"reason": "top 1.1% by frequency,growth=100%,not yet covered by interest keywords or stopwords. Evidence: total_count=8, days_seen=8, recent_count=8.",
|
||||
"suggestions_path": "outputs/term_index/review/term-cleanup-suggestions-2026-07-15.json",
|
||||
"suggestion_date": "2026-07-15",
|
||||
"based_on_days": 365
|
||||
},
|
||||
{
|
||||
"applied_at": "2026-07-15T02:17:50.155586Z",
|
||||
"action": "add_interest_keyword",
|
||||
"term": "CLI",
|
||||
"reason": "top 1.1% by frequency,growth=100%,not yet covered by interest keywords or stopwords. Evidence: total_count=8, days_seen=7, recent_count=8.",
|
||||
"suggestions_path": "outputs/term_index/review/term-cleanup-suggestions-2026-07-15.json",
|
||||
"suggestion_date": "2026-07-15",
|
||||
"based_on_days": 365
|
||||
},
|
||||
{
|
||||
"applied_at": "2026-07-15T02:17:50.155586Z",
|
||||
"action": "add_interest_keyword",
|
||||
"term": "子Agent",
|
||||
"reason": "top 1.3% by frequency,growth=100%,not yet covered by interest keywords or stopwords. Evidence: total_count=7, days_seen=6, recent_count=7.",
|
||||
"suggestions_path": "outputs/term_index/review/term-cleanup-suggestions-2026-07-15.json",
|
||||
"suggestion_date": "2026-07-15",
|
||||
"based_on_days": 365
|
||||
},
|
||||
{
|
||||
"applied_at": "2026-07-15T02:17:50.155586Z",
|
||||
"action": "add_interest_keyword",
|
||||
"term": "Loop Engineering",
|
||||
"reason": "top 1.5% by frequency,growth=100%,not yet covered by interest keywords or stopwords. Evidence: total_count=6, days_seen=6, recent_count=6.",
|
||||
"suggestions_path": "outputs/term_index/review/term-cleanup-suggestions-2026-07-15.json",
|
||||
"suggestion_date": "2026-07-15",
|
||||
"based_on_days": 365
|
||||
},
|
||||
{
|
||||
"applied_at": "2026-07-15T02:17:50.155586Z",
|
||||
"action": "add_interest_keyword",
|
||||
"term": "MoE",
|
||||
"reason": "top 2.0% by frequency,growth=100%,not yet covered by interest keywords or stopwords. Evidence: total_count=5, days_seen=4, recent_count=5.",
|
||||
"suggestions_path": "outputs/term_index/review/term-cleanup-suggestions-2026-07-15.json",
|
||||
"suggestion_date": "2026-07-15",
|
||||
"based_on_days": 365
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
+110
-3
@@ -1,19 +1,126 @@
|
||||
[
|
||||
"1688",
|
||||
"AGI",
|
||||
"AIHOT",
|
||||
"AI日报",
|
||||
"All In Code",
|
||||
"Andrej Karpathy",
|
||||
"Anthropic",
|
||||
"auto-twitter-campaign",
|
||||
"Boundaries",
|
||||
"ChangeSet",
|
||||
"Channels",
|
||||
"Claude Fable 5",
|
||||
"Claude Mythos",
|
||||
"Cohere",
|
||||
"Confidence Head",
|
||||
"Cosmos 3",
|
||||
"Databricks",
|
||||
"DINOv2",
|
||||
"domain-mapping",
|
||||
"Dropbox",
|
||||
"EchoGen",
|
||||
"FLUX.1-dev VAE",
|
||||
"GB300 GPU",
|
||||
"GLM 5.2",
|
||||
"GLM5.0",
|
||||
"GPT-5.5",
|
||||
"GPT-Live",
|
||||
"GPT5.5",
|
||||
"Grok 4.5",
|
||||
"GrowBrain",
|
||||
"iMedImage",
|
||||
"iMedLoop",
|
||||
"iMedMaaS",
|
||||
"iMedStudio",
|
||||
"J-space",
|
||||
"JLens",
|
||||
"John Jumper",
|
||||
"J空间",
|
||||
"KAIROS",
|
||||
"KubeRay",
|
||||
"LibTV Agent",
|
||||
"LingBot-Video",
|
||||
"Lumina",
|
||||
"Memory",
|
||||
"Markdown",
|
||||
"Marvis",
|
||||
"MDASH",
|
||||
"Meta Superintelligence Labs",
|
||||
"MTS",
|
||||
"Muse Image",
|
||||
"Muse Video",
|
||||
"N-gram Embedding",
|
||||
"OCP China",
|
||||
"OCP China 2026",
|
||||
"On-Policy Distillation",
|
||||
"OPC训练营",
|
||||
"OpenAI",
|
||||
"OpenBMC",
|
||||
"OpenClaw",
|
||||
"OpenViking",
|
||||
"Prompt",
|
||||
"Opus 4.8",
|
||||
"Qwen3",
|
||||
"Qwen3-30B-A3B",
|
||||
"RAS API",
|
||||
"Redfish",
|
||||
"ScMoE",
|
||||
"Seal AI Zone",
|
||||
"SealRouter",
|
||||
"Seedance 2.0",
|
||||
"Sonnet 5",
|
||||
"Spec模式",
|
||||
"STE固件团队",
|
||||
"Three.js",
|
||||
"Unity AI Gateway",
|
||||
"Vant Weapp",
|
||||
"WeTV",
|
||||
"WorkBuddy",
|
||||
"wpc",
|
||||
"YOLO Classifier",
|
||||
"一人公司",
|
||||
"中国科学技术大学",
|
||||
"五大扶持体系",
|
||||
"出门问问",
|
||||
"分镜",
|
||||
"剧本",
|
||||
"奋斗文化",
|
||||
"字节跳动",
|
||||
"小银",
|
||||
"得力",
|
||||
"德适科技",
|
||||
"成都天府长岛",
|
||||
"扣子",
|
||||
"星云平台",
|
||||
"火山引擎",
|
||||
"百度百舸",
|
||||
"百炼网关",
|
||||
"科大讯飞",
|
||||
"腾讯云开发者社区",
|
||||
"腾讯混元Hy3",
|
||||
"蚂蚁灵波",
|
||||
"贝尔实验室",
|
||||
"质量门禁",
|
||||
"配乐",
|
||||
"配音",
|
||||
"银行客户经理",
|
||||
"飞盘物理"
|
||||
"飞书妙搭",
|
||||
"飞盘物理",
|
||||
"自动化",
|
||||
"定时任务",
|
||||
"开源模型",
|
||||
"陌生化",
|
||||
"AlphaFold",
|
||||
"Brand Kit",
|
||||
"DataWorks",
|
||||
"Enhance-Nanocodec",
|
||||
"IRIS Codec",
|
||||
"Lovart",
|
||||
"MiniMax M3",
|
||||
"Gemini 3.5 Flash",
|
||||
"Codex",
|
||||
"CodeBuddy",
|
||||
"Claude Cowork",
|
||||
"AGENTS.md",
|
||||
"Claude",
|
||||
"RLVR"
|
||||
]
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"schema_version": "v1",
|
||||
"updated_at": "2026-05-14T09:12:50.749877Z",
|
||||
"updated_at": "2026-07-15T02:17:50.155586Z",
|
||||
"terms": [
|
||||
{
|
||||
"term": "A2A",
|
||||
|
||||
@@ -6,8 +6,30 @@ from .freshrss_pipeline_jobs import (
|
||||
start_freshrss_pipeline_job,
|
||||
)
|
||||
from .query_service import get_delivery_payload, get_run_report, get_run_status, list_run_artifacts, list_runs
|
||||
from .resume_jobs import get_resume_job_result, get_resume_job_status, start_resume_job
|
||||
from .resume_service import inspect_resume_plan, resume_run
|
||||
|
||||
# NOTE: resume_jobs and resume_service are NOT eagerly imported here to avoid
|
||||
# a circular import chain:
|
||||
# workflows/freshrss_pipeline.py -> runtime -> resume_jobs -> resume_service
|
||||
# -> workflows/freshrss_pipeline.py (circular!)
|
||||
# They are lazy-loaded via __getattr__ when accessed as summary_mcp.runtime.*
|
||||
|
||||
|
||||
def __getattr__(name):
|
||||
import importlib
|
||||
|
||||
_LAZY = {
|
||||
"get_resume_job_result": ("resume_jobs", "get_resume_job_result"),
|
||||
"get_resume_job_status": ("resume_jobs", "get_resume_job_status"),
|
||||
"start_resume_job": ("resume_jobs", "start_resume_job"),
|
||||
"inspect_resume_plan": ("resume_service", "inspect_resume_plan"),
|
||||
"resume_run": ("resume_service", "resume_run"),
|
||||
}
|
||||
if name in _LAZY:
|
||||
mod_name, attr_name = _LAZY[name]
|
||||
mod = importlib.import_module(f".{mod_name}", __package__)
|
||||
return getattr(mod, attr_name)
|
||||
raise AttributeError(f"module {__name__!r} has no attribute {name!r}")
|
||||
|
||||
|
||||
__all__ = [
|
||||
"ArtifactRecord",
|
||||
|
||||
@@ -9,6 +9,8 @@ from pathlib import Path
|
||||
from typing import Any
|
||||
from uuid import uuid4
|
||||
|
||||
from dotenv import dotenv_values
|
||||
|
||||
from .run_store import RunStore
|
||||
from .query_service import _resolve_run_record
|
||||
|
||||
@@ -29,6 +31,25 @@ DEFAULT_STAGES = [
|
||||
]
|
||||
MIN_JOB_STALE_SECONDS = 30 * 60
|
||||
MAX_JOB_STALE_SECONDS = 6 * 60 * 60
|
||||
DEFAULT_DOTENV_PATH = REPO_ROOT / ".env"
|
||||
|
||||
|
||||
def _build_subprocess_env() -> dict[str, str]:
|
||||
"""Build an env dict for subprocess, merging parent env with .env values.
|
||||
|
||||
The subprocess inherits the Hermes MCP server's environment, but .env values
|
||||
may not be in os.environ at the time the subprocess is spawned. This function
|
||||
loads them from .env and merges them in so the child process sees all needed
|
||||
variables (LLM_API_KEY, LLM_MODEL, LLM_API_URL, FRESHRSS_*, etc.) directly
|
||||
in os.environ, avoiding any dotenv-loading timing issues inside the subprocess.
|
||||
"""
|
||||
env = os.environ.copy()
|
||||
if DEFAULT_DOTENV_PATH.exists():
|
||||
for key, value in dotenv_values(DEFAULT_DOTENV_PATH).items():
|
||||
if isinstance(key, str) and isinstance(value, str) and value:
|
||||
# Only set if not already present in parent env
|
||||
env.setdefault(key, value)
|
||||
return env
|
||||
|
||||
|
||||
def _now() -> datetime:
|
||||
@@ -218,6 +239,7 @@ def start_freshrss_pipeline_job(
|
||||
stdout=subprocess.DEVNULL,
|
||||
stderr=subprocess.DEVNULL,
|
||||
start_new_session=True,
|
||||
env=_build_subprocess_env(),
|
||||
)
|
||||
except Exception as exc:
|
||||
report_file = _write_job_report(
|
||||
|
||||
@@ -9,6 +9,8 @@ from pathlib import Path
|
||||
from typing import Any
|
||||
from uuid import uuid4
|
||||
|
||||
from dotenv import dotenv_values
|
||||
|
||||
from .query_service import _resolve_run_record
|
||||
from .resume_service import (
|
||||
SUPPORTED_RESUME_STAGES,
|
||||
@@ -38,6 +40,22 @@ DEFAULT_STAGES = [
|
||||
]
|
||||
MIN_JOB_STALE_SECONDS = 30 * 60
|
||||
MAX_JOB_STALE_SECONDS = 6 * 60 * 60
|
||||
DEFAULT_DOTENV_PATH = REPO_ROOT / ".env"
|
||||
|
||||
|
||||
def _build_subprocess_env() -> dict[str, str]:
|
||||
"""Build an env dict for subprocess, merging parent env with .env values.
|
||||
|
||||
Ensures the subprocess sees all needed variables (LLM_API_KEY, LLM_MODEL,
|
||||
LLM_API_URL, FRESHRSS_*, etc.) directly in os.environ, avoiding dotenv
|
||||
timing issues in the child process.
|
||||
"""
|
||||
env = os.environ.copy()
|
||||
if DEFAULT_DOTENV_PATH.exists():
|
||||
for key, value in dotenv_values(DEFAULT_DOTENV_PATH).items():
|
||||
if isinstance(key, str) and isinstance(value, str) and value:
|
||||
env.setdefault(key, value)
|
||||
return env
|
||||
|
||||
|
||||
def _now() -> datetime:
|
||||
@@ -276,6 +294,7 @@ def start_resume_job(*, run_id: str) -> dict[str, Any]:
|
||||
stdout=subprocess.DEVNULL,
|
||||
stderr=subprocess.DEVNULL,
|
||||
start_new_session=True,
|
||||
env=_build_subprocess_env(),
|
||||
)
|
||||
except Exception as exc:
|
||||
report_file = _write_job_report(
|
||||
|
||||
@@ -2,6 +2,7 @@ from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
from concurrent.futures import ThreadPoolExecutor, as_completed
|
||||
from datetime import date, datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
@@ -192,7 +193,7 @@ def _build_candidate_batch_payload(*, run_id: str, item_contexts: list[dict[str,
|
||||
}
|
||||
|
||||
|
||||
def _persist_summary_batch_artifact(*, run_store: RunStore, run_dir: Path, item_contexts: list[dict[str, Any]]) -> Path:
|
||||
def _persist_summary_batch_artifact(*, run_store, run_dir: Path, item_contexts: list[dict[str, Any]]) -> Path:
|
||||
output_path = _summary_batch_output(run_dir)
|
||||
_save_json(
|
||||
output_path,
|
||||
@@ -202,7 +203,7 @@ def _persist_summary_batch_artifact(*, run_store: RunStore, run_dir: Path, item_
|
||||
return output_path
|
||||
|
||||
|
||||
def _persist_candidate_batch_artifact(*, run_store: RunStore, run_dir: Path, item_contexts: list[dict[str, Any]]) -> Path:
|
||||
def _persist_candidate_batch_artifact(*, run_store, run_dir: Path, item_contexts: list[dict[str, Any]]) -> Path:
|
||||
output_path = _candidate_batch_output(run_dir)
|
||||
_save_json(
|
||||
output_path,
|
||||
@@ -461,37 +462,50 @@ def run_freshrss_pipeline(
|
||||
summary_success_count = 0
|
||||
summary_failed_count = 0
|
||||
summary_candidates = [ctx for ctx in item_contexts if ctx["extraction"] is not None and ctx["extraction"].success]
|
||||
for item_context in summary_candidates:
|
||||
item_report = item_context["item_report"]
|
||||
summary_exit_code, summary_payload, summary_report = run_loop_payload(
|
||||
extracted_payload=item_context["extracted_payload"],
|
||||
prompt_path=resolved_prompt_path,
|
||||
output_path=item_context["summary_output"],
|
||||
max_retries=max_retries,
|
||||
timeout_seconds=timeout_seconds,
|
||||
api_key=resolved_llm_api_key,
|
||||
model=resolved_llm_model,
|
||||
api_url=resolved_llm_api_url,
|
||||
)
|
||||
if summary_exit_code != 0 or summary_payload is None:
|
||||
item_report["status"] = "summary_failed"
|
||||
if summary_report is not None:
|
||||
item_report["summary_errors"] = summary_report.errors
|
||||
summary_failed_count += 1
|
||||
else:
|
||||
item_context["summary_payload"] = summary_payload
|
||||
item_report["status"] = "summarized"
|
||||
summary_success_count += 1
|
||||
# Parallelize LLM summaries — I/O bound calls, independent per article
|
||||
with ThreadPoolExecutor(max_workers=min(len(summary_candidates) or 1, 4)) as pool:
|
||||
fut_map = {}
|
||||
for item_context in summary_candidates:
|
||||
fut = pool.submit(
|
||||
run_loop_payload,
|
||||
extracted_payload=item_context["extracted_payload"],
|
||||
prompt_path=resolved_prompt_path,
|
||||
output_path=item_context["summary_output"],
|
||||
max_retries=max_retries,
|
||||
timeout_seconds=timeout_seconds,
|
||||
api_key=resolved_llm_api_key,
|
||||
model=resolved_llm_model,
|
||||
api_url=resolved_llm_api_url,
|
||||
)
|
||||
fut_map[fut] = item_context
|
||||
|
||||
run_store.update_stage(
|
||||
SUMMARY_STAGE,
|
||||
outputs={
|
||||
"expected_items": extracted_success_count,
|
||||
"completed_items": summary_success_count + summary_failed_count,
|
||||
"success_count": summary_success_count,
|
||||
"failed_count": summary_failed_count,
|
||||
},
|
||||
)
|
||||
for fut in as_completed(fut_map):
|
||||
item_context = fut_map[fut]
|
||||
item_report = item_context["item_report"]
|
||||
try:
|
||||
summary_exit_code, summary_payload, summary_report = fut.result()
|
||||
except Exception as exc:
|
||||
summary_exit_code, summary_payload, summary_report = 1, None, None
|
||||
|
||||
if summary_exit_code != 0 or summary_payload is None:
|
||||
item_report["status"] = "summary_failed"
|
||||
if summary_report is not None:
|
||||
item_report["summary_errors"] = summary_report.errors
|
||||
summary_failed_count += 1
|
||||
else:
|
||||
item_context["summary_payload"] = summary_payload
|
||||
item_report["status"] = "summarized"
|
||||
summary_success_count += 1
|
||||
|
||||
run_store.update_stage(
|
||||
SUMMARY_STAGE,
|
||||
outputs={
|
||||
"expected_items": extracted_success_count,
|
||||
"completed_items": summary_success_count + summary_failed_count,
|
||||
"success_count": summary_success_count,
|
||||
"failed_count": summary_failed_count,
|
||||
},
|
||||
)
|
||||
|
||||
summary_batch_output = _persist_summary_batch_artifact(
|
||||
run_store=run_store,
|
||||
|
||||
Reference in New Issue
Block a user