diff --git a/CODE_REVIEW.md b/CODE_REVIEW.md index 2c0f69a..9b6cd2f 100644 --- a/CODE_REVIEW.md +++ b/CODE_REVIEW.md @@ -8,7 +8,7 @@ ## 重构建议(中等成本) -- [ ] `run_freshrss_pipeline` 函数过长(约300行)— 拆分为 `_process_single_item()`、`_build_and_persist_delivery()` 等内部函数,主函数只做编排 +- [x] `run_freshrss_pipeline` 函数过长(约300行)— 拆分为 `_process_item()`、`_build_and_persist_delivery()`、`_build_run_report()` 三个内部函数,主函数只做编排(已完成) - [ ] `load_filter_rules` 每次 pipeline 调用都重新读文件 — 加模块级缓存,MCP 服务长期运行时避免重复 I/O ## 功能补全 diff --git a/refactor-plan.md b/refactor-plan.md new file mode 100644 index 0000000..b856374 --- /dev/null +++ b/refactor-plan.md @@ -0,0 +1,54 @@ +# 重构计划:拆分 run_freshrss_pipeline 为内部辅助函数 + +## 背景 + +`src/summary_mcp/workflows/freshrss_pipeline.py` 中的 `run_freshrss_pipeline` 函数约 330 行, +将 6 个阶段全部写在一个函数体内,阅读、测试和后续扩展(如并发、重试策略)都比较困难。 +目标是在不改变任何外部行为的前提下,提取 3 个内部辅助函数。 + +当前无测试覆盖,验证方式为函数签名和返回值结构保持不变。 + +## 涉及文件 + +- `src/summary_mcp/workflows/freshrss_pipeline.py`(唯一修改文件) + +## 提取 3 个内部辅助函数 + +### 1. `_process_item(...)` — 单条 item 处理(当前 160-258 行) + +提取 for 循环体(约 100 行)为独立函数。 +返回 `item_report` dict;当 status 为 `delivered` 时,额外携带 `_candidate` 和 `_external_id` +两个临时键供调用方解包,写盘前剥离这两个键。 +用 `return item_report` 替代循环中的 `continue`。 + +### 2. `_build_and_persist_delivery(...)` — 阶段 4+5(当前 260-279 行) + +提取 payload 构建 + 词元索引持久化。 +返回 `(delivery_payload, keyword_index_result)`。 + +### 3. `_build_run_report(...)` — 报告组装(当前 288-308 行) + +提取 status_counts 统计 + report dict 构建。 +返回 report dict,主函数拿到后再调用 `_save_json` 写盘。 + +## 重构后主函数结构(约 80 行) + +1. 阶段 1:初始化(不变) +2. 阶段 2:拉取 FreshRSS + 加载规则/context(不变) +3. 阶段 3:for 循环调用 `_process_item(...)`,从返回值解包 candidate +4. 阶段 4+5:`delivery_payload, keyword_index_result = _build_and_persist_delivery(...)` +5. 阶段 6:标记已读,`report = _build_run_report(...)`,写盘,返回 + +## 约束 + +- `run_freshrss_pipeline` 外部签名不变 +- 返回 dict 的键结构不变 +- 所有文件写入路径不变 +- 纯结构性重构,无行为变化 +- 无需新增 import + +## 验证 + +重构完成后运行: + + python -c "from summary_mcp.workflows.freshrss_pipeline import run_freshrss_pipeline; print('ok')" diff --git a/src/summary_mcp/workflows/freshrss_pipeline.py b/src/summary_mcp/workflows/freshrss_pipeline.py index eb2bb2b..08088a6 100644 --- a/src/summary_mcp/workflows/freshrss_pipeline.py +++ b/src/summary_mcp/workflows/freshrss_pipeline.py @@ -71,6 +71,187 @@ def _maybe_path(enabled: bool, path: Path) -> Path | None: return path if enabled else None +def _process_item( + *, + index: int, + item: Any, + resolved_output_dir: Path, + resolved_prompt_path: Path, + resolved_run_id: str, + debug_artifacts: bool, + loaded_rules: list, + filter_context: Any, + max_retries: int, + timeout_seconds: float, + resolved_llm_api_key: str, + resolved_llm_model: str, + resolved_llm_api_url: str, +) -> dict[str, Any]: + # 处理单条 item:提取 -> LLM 摘要 -> 规则过滤 -> 候选构建 + # 返回 item_report dict;delivered 时额外携带 _candidate/_external_id 供调用方解包 + item_key = f"item-{index:02d}" + item_path = _maybe_path(debug_artifacts, resolved_output_dir / "items" / f"{item_key}.item.json") + extracted_path = resolved_output_dir / "extracted" / f"{item_key}.extracted.json" + summary_output = _maybe_path(debug_artifacts, resolved_output_dir / "summary" / item_key / "result.loop.json") + filter_path = _maybe_path(debug_artifacts, resolved_output_dir / "filter" / f"{item_key}.filter.json") + record_path = _maybe_path(debug_artifacts, resolved_output_dir / "candidates" / f"{item_key}.article-candidate-record.json") + openclaw_path = _maybe_path(debug_artifacts, resolved_output_dir / "candidates" / f"{item_key}.openclaw-candidate-input.json") + + if item_path is not None: + _save_json(item_path, item.model_dump(mode="json")) + + item_report: dict[str, Any] = { + "item_key": item_key, + "item_id": item.item_id, + "external_id": item.external_id, + "url": str(item.url), + "title": item.title, + "status": "pulled", + } + if debug_artifacts: + item_report["paths"] = { + "item": str(item_path) if item_path else None, + "extracted": str(extracted_path) if extracted_path else None, + "summary": str(summary_output) if summary_output else None, + "filter": str(filter_path) if filter_path else None, + "article_candidate": str(record_path) if record_path else None, + "openclaw_candidate": str(openclaw_path) if openclaw_path else None, + } + + extraction = extract_content(ExtractionInput(item=item)) + extracted_payload = extraction.model_dump(mode="json") + if extracted_path is not None: + _save_json(extracted_path, extracted_payload) + if not extraction.success or extraction.article is None: + item_report["status"] = "extract_failed" + item_report["error"] = extraction.error.model_dump(mode="json") if extraction.error else None + return item_report + + item_report["status"] = "extracted" + summary_exit_code, summary_payload, summary_report = run_loop_payload( + extracted_payload=extracted_payload, + prompt_path=resolved_prompt_path, + output_path=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 + return item_report + + summary = LlmSummaryResult.model_validate(summary_payload) + decision = evaluate_filter_rules( + FilterInput(item=item, article=extraction.article, summary=summary, context=filter_context), + loaded_rules, + ) + if filter_path is not None: + _save_json(filter_path, decision.model_dump(mode="json")) + + record = build_article_candidate_record( + summary=summary, + article=extraction.article, + filter_result=decision, + item=item, + source_refs=CandidateSourceRefs( + item_path=str(item_path) if item_path else None, + extracted_path=str(extracted_path) if extracted_path else None, + summary_path=str(summary_output) if summary_output else None, + filter_path=str(filter_path) if filter_path else None, + ), + metadata=CandidateMetadata( + generated_at=datetime.now(tz=UTC), + producer="run_freshrss_pipeline", + run_id=resolved_run_id, + ), + ) + openclaw_input = build_openclaw_candidate_input(record) + + if record_path is not None: + _save_json(record_path, record.model_dump(mode="json")) + if openclaw_path is not None: + _save_json(openclaw_path, openclaw_input.model_dump(mode="json")) + + item_report["status"] = "delivered" + item_report["selection_decision"] = decision.decision + item_report["candidate_id"] = openclaw_input.candidate_id + item_report["_candidate"] = openclaw_input + item_report["_external_id"] = item.external_id + return item_report + + +def _build_and_persist_delivery( + *, + delivered_candidates: list[OpenClawCandidateInput], + resolved_run_id: str, + resolved_delivery_date: Any, + delivery_output: Path, +) -> tuple[OpenClawDeliveryPayload, dict[str, Any]]: + # 按 digest_rank 排序,构建 delivery payload,持久化词元索引 + delivered_candidates.sort(key=lambda candidate: candidate.digest_rank, reverse=True) + delivery_payload = build_openclaw_delivery_payload( + delivered_candidates, + run_id=resolved_run_id, + for_date=resolved_delivery_date, + ) + _save_json(delivery_output, delivery_payload.model_dump(mode="json")) + + keyword_index_result = persist_keyword_indexes( + delivery_payload.candidates, + for_date=delivery_payload.date, + digest_id=delivery_payload.run_id, + source="openclaw_delivery_payload", + daily_dir=DEFAULT_TERM_DAILY_DIR, + stats_path=DEFAULT_TERM_STATS_PATH, + aliases_path=DEFAULT_TERM_ALIASES_PATH, + stopwords_path=DEFAULT_TERM_STOPWORDS_PATH, + ) + return delivery_payload, keyword_index_result + + +def _build_run_report( + *, + resolved_run_id: str, + started_at: datetime, + limit: int, + items: list, + delivered_candidates: list, + marked_count: int, + mark_read: bool, + debug_artifacts: bool, + raw_output: Path, + delivery_output: Path, + keyword_index_result: dict[str, Any], + item_reports: list[dict[str, Any]], +) -> dict[str, Any]: + # 统计各状态计数,组装 run report dict + status_counts: dict[str, int] = {} + for item_report in item_reports: + status = str(item_report["status"]) + status_counts[status] = status_counts.get(status, 0) + 1 + + return { + "run_id": resolved_run_id, + "started_at": started_at.isoformat(), + "completed_at": datetime.now(tz=UTC).isoformat(), + "requested_limit": limit, + "pulled_count": len(items), + "delivered_count": len(delivered_candidates), + "marked_read_count": marked_count, + "mark_read_requested": mark_read, + "debug_artifacts": debug_artifacts, + "raw_output": str(raw_output), + "delivery_output": str(delivery_output), + "keyword_index": keyword_index_result, + "status_counts": status_counts, + "items": item_reports, + } + + def run_freshrss_pipeline( *, api_base_url: str | None = None, @@ -95,6 +276,7 @@ def run_freshrss_pipeline( delivery_date: date | None = None, output_dir: Path | None = None, ) -> dict[str, Any]: + # --- 阶段 1:初始化 run_id、输出路径、凭证 --- started_at = datetime.now(tz=UTC) resolved_output_dir = output_dir or default_output_dir() run_stamp = started_at.strftime("%Y%m%d-%H%M%S") @@ -117,6 +299,7 @@ def run_freshrss_pipeline( report_output = resolved_output_dir / "run-report.json" items_list_output = _maybe_path(debug_artifacts, resolved_output_dir / "items" / "freshrss.items.json") + # --- 阶段 2:登录 FreshRSS,拉取未读条目,写原始 payload --- client = FreshRSSClient( api_base_url=resolved_api_base_url, username=resolved_username, @@ -150,156 +333,63 @@ def run_freshrss_pipeline( else: filter_context = FilterContext() + # --- 阶段 3:逐条处理(提取 -> LLM 摘要 -> 规则过滤 -> 候选构建) --- delivered_candidates: list[OpenClawCandidateInput] = [] delivered_item_ids: list[str] = [] item_reports: list[dict[str, Any]] = [] for index, item in enumerate(items, start=1): - # 逐条处理:提取 -> 摘要 -> 过滤 -> 候选构建;任一步骤失败则记录状态后 continue - item_key = f"item-{index:02d}" - item_path = _maybe_path(debug_artifacts, resolved_output_dir / "items" / f"{item_key}.item.json") - extracted_path = resolved_output_dir / "extracted" / f"{item_key}.extracted.json" - summary_output = _maybe_path(debug_artifacts, resolved_output_dir / "summary" / item_key / "result.loop.json") - filter_path = _maybe_path(debug_artifacts, resolved_output_dir / "filter" / f"{item_key}.filter.json") - record_path = _maybe_path(debug_artifacts, resolved_output_dir / "candidates" / f"{item_key}.article-candidate-record.json") - openclaw_path = _maybe_path(debug_artifacts, resolved_output_dir / "candidates" / f"{item_key}.openclaw-candidate-input.json") - - if item_path is not None: - _save_json(item_path, item.model_dump(mode="json")) - - item_report: dict[str, Any] = { - "item_key": item_key, - "item_id": item.item_id, - "external_id": item.external_id, - "url": str(item.url), - "title": item.title, - "status": "pulled", - } - if debug_artifacts: - item_report["paths"] = { - "item": str(item_path) if item_path else None, - "extracted": str(extracted_path) if extracted_path else None, - "summary": str(summary_output) if summary_output else None, - "filter": str(filter_path) if filter_path else None, - "article_candidate": str(record_path) if record_path else None, - "openclaw_candidate": str(openclaw_path) if openclaw_path else None, - } - - extraction = extract_content(ExtractionInput(item=item)) - extracted_payload = extraction.model_dump(mode="json") - if extracted_path is not None: - _save_json(extracted_path, extracted_payload) - if not extraction.success or extraction.article is None: - item_report["status"] = "extract_failed" - item_report["error"] = extraction.error.model_dump(mode="json") if extraction.error else None - item_reports.append(item_report) - continue - - item_report["status"] = "extracted" - summary_exit_code, summary_payload, summary_report = run_loop_payload( - extracted_payload=extracted_payload, - prompt_path=resolved_prompt_path, - output_path=summary_output, + result = _process_item( + index=index, + item=item, + resolved_output_dir=resolved_output_dir, + resolved_prompt_path=resolved_prompt_path, + resolved_run_id=resolved_run_id, + debug_artifacts=debug_artifacts, + loaded_rules=loaded_rules, + filter_context=filter_context, max_retries=max_retries, timeout_seconds=timeout_seconds, - api_key=resolved_llm_api_key, - model=resolved_llm_model, - api_url=resolved_llm_api_url, + resolved_llm_api_key=resolved_llm_api_key, + resolved_llm_model=resolved_llm_model, + resolved_llm_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 - item_reports.append(item_report) - continue + if result.get("status") == "delivered": + delivered_candidates.append(result.pop("_candidate")) + external_id = result.pop("_external_id", None) + if external_id: + delivered_item_ids.append(external_id) + item_reports.append(result) - summary = LlmSummaryResult.model_validate(summary_payload) - decision = evaluate_filter_rules( - FilterInput(item=item, article=extraction.article, summary=summary, context=filter_context), - loaded_rules, - ) - if filter_path is not None: - _save_json(filter_path, decision.model_dump(mode="json")) - - record = build_article_candidate_record( - summary=summary, - article=extraction.article, - filter_result=decision, - item=item, - source_refs=CandidateSourceRefs( - item_path=str(item_path) if item_path else None, - extracted_path=str(extracted_path) if extracted_path else None, - summary_path=str(summary_output) if summary_output else None, - filter_path=str(filter_path) if filter_path else None, - ), - metadata=CandidateMetadata( - generated_at=datetime.now(tz=UTC), - producer="run_freshrss_pipeline", - run_id=resolved_run_id, - ), - ) - openclaw_input = build_openclaw_candidate_input(record) - - if record_path is not None: - _save_json(record_path, record.model_dump(mode="json")) - if openclaw_path is not None: - _save_json(openclaw_path, openclaw_input.model_dump(mode="json")) - - delivered_candidates.append(openclaw_input) - if item.external_id: - delivered_item_ids.append(item.external_id) - - item_report["status"] = "delivered" - item_report["selection_decision"] = decision.decision - item_report["candidate_id"] = openclaw_input.candidate_id - item_reports.append(item_report) - - delivered_candidates.sort(key=lambda candidate: candidate.digest_rank, reverse=True) - delivery_payload = build_openclaw_delivery_payload( - delivered_candidates, - run_id=resolved_run_id, - for_date=resolved_delivery_date, - ) - _save_json(delivery_output, delivery_payload.model_dump(mode="json")) - - keyword_index_result = persist_keyword_indexes( - delivery_payload.candidates, - for_date=delivery_payload.date, - digest_id=delivery_payload.run_id, - source="openclaw_delivery_payload", - daily_dir=DEFAULT_TERM_DAILY_DIR, - stats_path=DEFAULT_TERM_STATS_PATH, - aliases_path=DEFAULT_TERM_ALIASES_PATH, - stopwords_path=DEFAULT_TERM_STOPWORDS_PATH, + # --- 阶段 4+5:构建 delivery payload 并持久化词元索引 --- + delivery_payload, keyword_index_result = _build_and_persist_delivery( + delivered_candidates=delivered_candidates, + resolved_run_id=resolved_run_id, + resolved_delivery_date=resolved_delivery_date, + delivery_output=delivery_output, ) + # --- 阶段 6:标记已读,汇总报告,返回结果 --- marked_count = 0 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) marked_count = len({item_id for item_id in delivered_item_ids if item_id}) - status_counts: dict[str, int] = {} - for item_report in item_reports: - status = str(item_report["status"]) - status_counts[status] = status_counts.get(status, 0) + 1 - - report = { - "run_id": resolved_run_id, - "started_at": started_at.isoformat(), - "completed_at": datetime.now(tz=UTC).isoformat(), - "requested_limit": limit, - "pulled_count": len(items), - "delivered_count": len(delivered_candidates), - "marked_read_count": marked_count, - "mark_read_requested": mark_read, - "debug_artifacts": debug_artifacts, - "raw_output": str(raw_output), - "delivery_output": str(delivery_output), - "keyword_index": keyword_index_result, - "status_counts": status_counts, - "items": item_reports, - } + report = _build_run_report( + resolved_run_id=resolved_run_id, + started_at=started_at, + limit=limit, + items=items, + delivered_candidates=delivered_candidates, + marked_count=marked_count, + mark_read=mark_read, + debug_artifacts=debug_artifacts, + raw_output=raw_output, + delivery_output=delivery_output, + keyword_index_result=keyword_index_result, + item_reports=item_reports, + ) _save_json(report_output, report) return { @@ -312,7 +402,7 @@ def run_freshrss_pipeline( "pulled_count": len(items), "delivered_count": len(delivered_candidates), "marked_read_count": marked_count, - "status_counts": status_counts, + "status_counts": report["status_counts"], "debug_artifacts": debug_artifacts, "delivery_payload": delivery_payload.model_dump(mode="json"), "items": item_reports,