Add inline comments to core modules and code review notes

Added explanatory comments to server.py, freshrss_pipeline.py, pipeline.py,
freshrss.py, keyword_index.py, summary_loop.py, and filters/engine.py.
Also added CODE_REVIEW.md with prioritized improvement suggestions.
This commit is contained in:
wdm
2026-03-28 23:59:50 +08:00
parent 3940db1443
commit 41f7ac0e91
8 changed files with 56 additions and 0 deletions
+17
View File
@@ -0,0 +1,17 @@
# 代码审查建议(2026-03-28)
## 立即可改(低成本)
- [ ] `server.py` context 通过临时文件传递绕路 — `run_freshrss_pipeline` 改为同时接受 `dict | Path` 类型的 context,消除 `NamedTemporaryFile` 绕路
- [ ] `_load_required_env` 报错信息不区分"未传参数"还是"环境变量未设",改善调试体验
- [ ] `evaluate_filter_rules` 决策逻辑歧义 — `stop_on_match=True` 命中后应直接以该规则 decision 为最终结果,而非继续聚合所有 matched
## 重构建议(中等成本)
- [ ] `run_freshrss_pipeline` 函数过长(约300行)— 拆分为 `_process_single_item()`、`_build_and_persist_delivery()` 等内部函数,主函数只做编排
- [ ] `load_filter_rules` 每次 pipeline 调用都重新读文件 — 加模块级缓存,MCP 服务长期运行时避免重复 I/O
## 功能补全
- [ ] `mark_read` 只标记 delivered items,drop/review 的 item 下次仍会重复拉取 — 引入 `mark_read_all_processed` 选项,或在报告中明确标注
- [ ] LLM 调用逐条串行 — 考虑用 `asyncio` + `httpx.AsyncClient` 并发处理多条 item,减少整体延迟
+5
View File
@@ -1,5 +1,8 @@
from __future__ import annotations from __future__ import annotations
# 日报级关键词索引:从当日 delivery 候选中提取关键词,
# 应用别名归并和停用词过滤,输出 DailyKeywordIndex 并聚合全局 KeywordStatsIndex。
import json import json
from collections import Counter from collections import Counter
from datetime import UTC, date, datetime from datetime import UTC, date, datetime
@@ -83,6 +86,7 @@ def normalize_keyword(
aliases: dict[str, str], aliases: dict[str, str],
stopwords: set[str], stopwords: set[str],
) -> str | None: ) -> str | None:
# 对单个关键词做别名替换 + 停用词过滤,返回 None 表示丢弃
raw_value = keyword.strip() raw_value = keyword.strip()
if not raw_value: if not raw_value:
return None return None
@@ -200,6 +204,7 @@ def persist_keyword_indexes(
stopwords_path: Path | None = None, stopwords_path: Path | None = None,
include_decisions: set[str] | None = None, include_decisions: set[str] | None = None,
) -> dict[str, object]: ) -> dict[str, object]:
# 构建并写入当日词元索引,同时全量重建全局词频统计;每次 pipeline 运行后自动调用
aliases = load_term_aliases(aliases_path) aliases = load_term_aliases(aliases_path)
stopwords = load_term_stopwords(stopwords_path) stopwords = load_term_stopwords(stopwords_path)
daily_index = build_daily_keyword_index( daily_index = build_daily_keyword_index(
+5
View File
@@ -1,5 +1,8 @@
from __future__ import annotations from __future__ import annotations
# 内容提取主流程:根据 item 来源决定使用内联 RSS 内容还是回源抓取,
# 再经过文本提取、质量评估、文章构建,输出标准化 ExtractionOutput。
from summary_mcp.core.content_loader import choose_inline_content, fetch_html from summary_mcp.core.content_loader import choose_inline_content, fetch_html
from summary_mcp.core.errors import SummaryError from summary_mcp.core.errors import SummaryError
from summary_mcp.core.extractor import extract_plain_text, extract_title from summary_mcp.core.extractor import extract_plain_text, extract_title
@@ -9,10 +12,12 @@ from summary_mcp.core.quality_checker import assess_quality
from summary_mcp.models.summary_io import DebugInfo, ExtractionInput, ExtractionOutput from summary_mcp.models.summary_io import DebugInfo, ExtractionInput, ExtractionOutput
# FreshRSS 来源的 item 不回源抓网页,直接使用 RSS 内联内容
RSS_ONLY_UPSTREAMS = {"freshrss"} RSS_ONLY_UPSTREAMS = {"freshrss"}
def _should_skip_fetch(extraction_input: ExtractionInput) -> bool: def _should_skip_fetch(extraction_input: ExtractionInput) -> bool:
# 判断当前 item 是否属于 RSS-only 上游,若是则禁止回源抓取
item = extraction_input.item item = extraction_input.item
if item is None: if item is None:
return False return False
+3
View File
@@ -1,5 +1,8 @@
from __future__ import annotations from __future__ import annotations
# LLM 摘要循环:构造提示词 -> 调用 LLM -> 解析 JSON -> 校验 -> 失败时生成修复提示词重试。
# 核心入口:run_loop_payload(),最多重试 max_retries 次。
import json import json
import os import os
import re import re
+6
View File
@@ -1,5 +1,8 @@
from __future__ import annotations from __future__ import annotations
# 确定性规则引擎:从 JSON 配置加载规则,按优先级评估每条规则的条件,
# 输出 keep / review / drop 决策及命中规则列表。
import json import json
from pathlib import Path from pathlib import Path
from typing import Any from typing import Any
@@ -32,6 +35,7 @@ def _normalize_value(value: Any) -> Any:
def _resolve_field(filter_input: FilterInput, field_path: str) -> Any: def _resolve_field(filter_input: FilterInput, field_path: str) -> Any:
# 按点分路径从 FilterInput 中取值,路径不存在时返回 None
current: Any = filter_input current: Any = filter_input
for part in field_path.split("."): for part in field_path.split("."):
current = _normalize_value(current) current = _normalize_value(current)
@@ -94,6 +98,8 @@ def _rule_matches(filter_input: FilterInput, rule: FilterRule) -> bool:
def evaluate_filter_rules(filter_input: FilterInput, rules: list[FilterRule]) -> FilterDecisionResult: def evaluate_filter_rules(filter_input: FilterInput, rules: list[FilterRule]) -> FilterDecisionResult:
# 按优先级排序后逐条评估规则;stop_on_match=True 时命中即终止
# 最终决策优先级:drop > keep > review;无规则命中时默认 review
matched: list[MatchedRule] = [] matched: list[MatchedRule] = []
for rule in sorted(rules, key=lambda item: item.action.priority, reverse=True): for rule in sorted(rules, key=lambda item: item.action.priority, reverse=True):
+6
View File
@@ -1,5 +1,8 @@
from __future__ import annotations from __future__ import annotations
# FreshRSS GReader API 集成:负责登录鉴权、拉取未读条目、将 entry 映射为标准 item、
# 以及成功投递后将条目标记为已读。
import hashlib import hashlib
from datetime import UTC, datetime from datetime import UTC, datetime
from typing import Any from typing import Any
@@ -12,6 +15,7 @@ from summary_mcp.models.item import Item
READ_TAG = "user/-/state/com.google/read" READ_TAG = "user/-/state/com.google/read"
KEPT_UNREAD_TAG = "user/-/state/com.google/kept-unread" KEPT_UNREAD_TAG = "user/-/state/com.google/kept-unread"
# RSS 内联内容低于此长度时视为无效,item 将被标记为 pending 并跳过
MIN_INLINE_CONTENT_LENGTH = 500 MIN_INLINE_CONTENT_LENGTH = 500
@@ -124,6 +128,8 @@ def _has_usable_inline_content(value: str | None) -> bool:
def map_entry_to_item(entry: dict[str, Any]) -> Item: def map_entry_to_item(entry: dict[str, Any]) -> Item:
# 将 FreshRSS GReader API 返回的原始 entry dict 映射为标准化 Item 对象
# fetch_state: 若内联内容足够长则为 'fetched',否则为 'pending'
url = _pick_entry_url(entry) url = _pick_entry_url(entry)
summary_html = _pick_content_block(entry, "summary") summary_html = _pick_content_block(entry, "summary")
content_html = _pick_content_block(entry, "content") content_html = _pick_content_block(entry, "content")
+7
View File
@@ -1,5 +1,8 @@
from __future__ import annotations from __future__ import annotations
# MCP 服务入口:将内容提取、过滤、FreshRSS 全链路管道暴露为 MCP 工具。
# 生产主入口是 run_freshrss_openclaw_pipeline,其余工具供单步调试使用。
import json import json
import tempfile import tempfile
from datetime import date from datetime import date
@@ -22,6 +25,7 @@ mcp = FastMCP(name="content-extract-mcp")
@mcp.tool() @mcp.tool()
def extract_url_content(url: str, language_hint: str | None = None) -> dict: def extract_url_content(url: str, language_hint: str | None = None) -> dict:
# 从单个 URL 抓取并提取结构化文章内容,供单步调试使用
"""Extract structured article content from a single URL.""" """Extract structured article content from a single URL."""
result = extract_content( result = extract_content(
ExtractionInput( ExtractionInput(
@@ -34,6 +38,7 @@ def extract_url_content(url: str, language_hint: str | None = None) -> dict:
@mcp.tool() @mcp.tool()
def extract_item_content(item: dict) -> dict: def extract_item_content(item: dict) -> dict:
# 从已标准化的 item 对象提取内容(跳过网络抓取,使用 RSS 内联内容)
"""Extract structured article content from a normalized item object.""" """Extract structured article content from a normalized item object."""
parsed_item = Item.model_validate(item) parsed_item = Item.model_validate(item)
result = extract_content(ExtractionInput(item=parsed_item)) result = extract_content(ExtractionInput(item=parsed_item))
@@ -47,6 +52,7 @@ def filter_summary_result(
item: dict | None = None, item: dict | None = None,
context: dict | None = None, context: dict | None = None,
) -> dict: ) -> dict:
# 对结构化摘要结果运行确定性规则引擎,返回 keep/review/drop 决策
"""Apply deterministic filter rules to a structured summary result.""" """Apply deterministic filter rules to a structured summary result."""
parsed_summary = LlmSummaryResult.model_validate(summary_result) parsed_summary = LlmSummaryResult.model_validate(summary_result)
parsed_article = ExtractedArticle.model_validate(extracted_article) if extracted_article else None parsed_article = ExtractedArticle.model_validate(extracted_article) if extracted_article else None
@@ -88,6 +94,7 @@ def run_freshrss_openclaw_pipeline(
include_item_reports: bool = False, include_item_reports: bool = False,
) -> dict: ) -> dict:
"""Run the full FreshRSS -> extract -> LLM -> filter -> OpenClaw payload pipeline.""" """Run the full FreshRSS -> extract -> LLM -> filter -> OpenClaw payload pipeline."""
# context 若传入则序列化为临时文件再传给 pipeline(TODO: pipeline 改为直接接受 dict)
temp_context_path: Path | None = None temp_context_path: Path | None = None
try: try:
@@ -1,5 +1,9 @@
from __future__ import annotations from __future__ import annotations
# FreshRSS 全链路管道:拉取未读条目 -> 内容提取 -> LLM 摘要 -> 规则过滤 ->
# 构建 OpenClaw delivery payload -> 写盘 -> 词元统计 -> 标记已读。
# 生产入口:run_freshrss_pipeline(),由 MCP 工具 run_freshrss_openclaw_pipeline 调用。
import json import json
import os import os
from datetime import UTC, date, datetime from datetime import UTC, date, datetime
@@ -46,6 +50,7 @@ def _load_json(path: Path) -> dict[str, Any]:
def _load_required_env(name: str, value: str | None) -> str: def _load_required_env(name: str, value: str | None) -> str:
# 优先使用传入的 value,否则读取同名环境变量;两者均缺失时抛出 RuntimeError
resolved = value or os.environ.get(name) resolved = value or os.environ.get(name)
if not resolved: if not resolved:
raise RuntimeError(f"Missing required value: {name}.") raise RuntimeError(f"Missing required value: {name}.")
@@ -137,6 +142,7 @@ def run_freshrss_pipeline(
item_reports: list[dict[str, Any]] = [] item_reports: list[dict[str, Any]] = []
for index, item in enumerate(items, start=1): for index, item in enumerate(items, start=1):
# 逐条处理:提取 -> 摘要 -> 过滤 -> 候选构建;任一步骤失败则记录状态后 continue
item_key = f"item-{index:02d}" item_key = f"item-{index:02d}"
item_path = _maybe_path(debug_artifacts, resolved_output_dir / "items" / f"{item_key}.item.json") item_path = _maybe_path(debug_artifacts, resolved_output_dir / "items" / f"{item_key}.item.json")
extracted_path = _maybe_path(debug_artifacts, resolved_output_dir / "extracted" / f"{item_key}.extracted.json") extracted_path = _maybe_path(debug_artifacts, resolved_output_dir / "extracted" / f"{item_key}.extracted.json")
@@ -256,6 +262,7 @@ def run_freshrss_pipeline(
marked_count = 0 marked_count = 0
if mark_read and delivered_item_ids: if mark_read and delivered_item_ids:
# 仅标记成功投递(delivered)的 item;drop/review 的 item 保持未读状态
client.mark_items_as_read(auth_token=auth_token, item_ids=delivered_item_ids) client.mark_items_as_read(auth_token=auth_token, item_ids=delivered_item_ids)
marked_count = len({item_id for item_id in delivered_item_ids if item_id}) marked_count = len({item_id for item_id in delivered_item_ids if item_id})