Compare commits
6
Commits
7bb9bef12f
...
79509fd1ba
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
79509fd1ba | ||
|
|
3d58408687 | ||
|
|
78039eff47 | ||
|
|
a2f6f8a28d | ||
|
|
2d9a38e465 | ||
|
|
ced6142b59 |
@@ -1,5 +1,6 @@
|
|||||||
.claude/
|
.claude/
|
||||||
.codex/
|
.codex/
|
||||||
|
.idea/
|
||||||
__pycache__/
|
__pycache__/
|
||||||
*.pyc
|
*.pyc
|
||||||
*.egg-info/
|
*.egg-info/
|
||||||
|
|||||||
+3
-3
@@ -2,13 +2,13 @@
|
|||||||
|
|
||||||
## 立即可改(低成本)
|
## 立即可改(低成本)
|
||||||
|
|
||||||
- [ ] `server.py` context 通过临时文件传递绕路 — `run_freshrss_pipeline` 改为同时接受 `dict | Path` 类型的 context,消除 `NamedTemporaryFile` 绕路
|
- [x] `server.py` context 通过临时文件传递绕路 — `run_freshrss_pipeline` 改为同时接受 `dict | Path` 类型的 context,消除 `NamedTemporaryFile` 绕路(已修复)
|
||||||
- [ ] `_load_required_env` 报错信息不区分"未传参数"还是"环境变量未设",改善调试体验
|
- [ ] `_load_required_env` 报错信息不区分"未传参数"还是"环境变量未设",改善调试体验
|
||||||
- [ ] `evaluate_filter_rules` 决策逻辑歧义 — `stop_on_match=True` 命中后应直接以该规则 decision 为最终结果,而非继续聚合所有 matched
|
- [x] `evaluate_filter_rules` 决策逻辑歧义 — `stop_on_match=True` 命中后应直接以该规则 decision 为最终结果,而非继续聚合所有 matched(已确认:当前规则集安全,在 filter-rule-engine-usage.md 补充了 drop 规则必须设 stop_on_match=true 的约束说明)
|
||||||
|
|
||||||
## 重构建议(中等成本)
|
## 重构建议(中等成本)
|
||||||
|
|
||||||
- [ ] `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
|
- [ ] `load_filter_rules` 每次 pipeline 调用都重新读文件 — 加模块级缓存,MCP 服务长期运行时避免重复 I/O
|
||||||
|
|
||||||
## 功能补全
|
## 功能补全
|
||||||
|
|||||||
@@ -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')"
|
||||||
+20
-34
@@ -3,8 +3,6 @@ from __future__ import annotations
|
|||||||
# MCP 服务入口:将内容提取、过滤、FreshRSS 全链路管道暴露为 MCP 工具。
|
# MCP 服务入口:将内容提取、过滤、FreshRSS 全链路管道暴露为 MCP 工具。
|
||||||
# 生产主入口是 run_freshrss_openclaw_pipeline,其余工具供单步调试使用。
|
# 生产主入口是 run_freshrss_openclaw_pipeline,其余工具供单步调试使用。
|
||||||
|
|
||||||
import json
|
|
||||||
import tempfile
|
|
||||||
from datetime import date
|
from datetime import date
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
@@ -95,38 +93,26 @@ 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)
|
result = run_freshrss_pipeline(
|
||||||
temp_context_path: Path | None = None
|
api_base_url=api_base_url,
|
||||||
|
username=username,
|
||||||
try:
|
api_password=api_password,
|
||||||
if context is not None:
|
stream_id=stream_id,
|
||||||
with tempfile.NamedTemporaryFile("w", encoding="utf-8", suffix=".json", delete=False) as handle:
|
limit=limit,
|
||||||
json.dump(context, handle, ensure_ascii=False, indent=2)
|
continuation=continuation,
|
||||||
temp_context_path = Path(handle.name)
|
include_read=include_read,
|
||||||
|
mark_read=mark_read,
|
||||||
result = run_freshrss_pipeline(
|
debug_artifacts=debug_artifacts,
|
||||||
api_base_url=api_base_url,
|
context=context,
|
||||||
username=username,
|
max_retries=max_retries,
|
||||||
api_password=api_password,
|
timeout_seconds=timeout_seconds,
|
||||||
stream_id=stream_id,
|
llm_api_key=llm_api_key,
|
||||||
limit=limit,
|
llm_model=llm_model,
|
||||||
continuation=continuation,
|
llm_api_url=llm_api_url,
|
||||||
include_read=include_read,
|
run_id=run_id,
|
||||||
mark_read=mark_read,
|
delivery_date=date.fromisoformat(date_value) if date_value else None,
|
||||||
debug_artifacts=debug_artifacts,
|
output_dir=Path(output_dir) if output_dir else None,
|
||||||
context_path=temp_context_path,
|
)
|
||||||
max_retries=max_retries,
|
|
||||||
timeout_seconds=timeout_seconds,
|
|
||||||
llm_api_key=llm_api_key,
|
|
||||||
llm_model=llm_model,
|
|
||||||
llm_api_url=llm_api_url,
|
|
||||||
run_id=run_id,
|
|
||||||
delivery_date=date.fromisoformat(date_value) if date_value else None,
|
|
||||||
output_dir=Path(output_dir) if output_dir else None,
|
|
||||||
)
|
|
||||||
finally:
|
|
||||||
if temp_context_path and temp_context_path.exists():
|
|
||||||
temp_context_path.unlink()
|
|
||||||
|
|
||||||
if not include_item_reports:
|
if not include_item_reports:
|
||||||
result = {key: value for key, value in result.items() if key != "items"}
|
result = {key: value for key, value in result.items() if key != "items"}
|
||||||
|
|||||||
@@ -53,10 +53,14 @@ 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
|
# 优先使用传入的 value,否则读取同名环境变量;两者均缺失时抛出 RuntimeError
|
||||||
resolved = value or os.environ.get(name)
|
if value:
|
||||||
if not resolved:
|
return value
|
||||||
raise RuntimeError(f"Missing required value: {name}.")
|
env_value = os.environ.get(name)
|
||||||
return resolved
|
if env_value:
|
||||||
|
return env_value
|
||||||
|
raise RuntimeError(
|
||||||
|
f"Missing required value '{name}': not passed as argument and not set as environment variable."
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def default_output_dir() -> Path:
|
def default_output_dir() -> Path:
|
||||||
@@ -67,6 +71,194 @@ def _maybe_path(enabled: bool, path: Path) -> Path | None:
|
|||||||
return path if enabled else 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 供调用方解包
|
||||||
|
# 步骤 1:确定各中间文件路径(debug_artifacts=False 时大部分路径为 None,不写盘)
|
||||||
|
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"))
|
||||||
|
|
||||||
|
# 步骤 2:初始化 item_report,记录基础元信息;debug 模式下附加各中间文件路径
|
||||||
|
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,
|
||||||
|
}
|
||||||
|
|
||||||
|
# 步骤 3:内容提取(RSS 内联内容 or 回源抓取);extracted.json 始终写盘
|
||||||
|
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
|
||||||
|
|
||||||
|
# 步骤 4:LLM 摘要循环,失败时最多重试 max_retries 次
|
||||||
|
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
|
||||||
|
|
||||||
|
# 步骤 5:规则引擎过滤,产出 keep/review/drop 决策及 digest_rank
|
||||||
|
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"))
|
||||||
|
|
||||||
|
# 步骤 6:构建 ArticleCandidateRecord 和 OpenClawCandidateInput
|
||||||
|
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"))
|
||||||
|
|
||||||
|
# 步骤 7:标记 delivered,附加临时键 _candidate/_external_id 供主函数解包后加入投递列表
|
||||||
|
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(
|
def run_freshrss_pipeline(
|
||||||
*,
|
*,
|
||||||
api_base_url: str | None = None,
|
api_base_url: str | None = None,
|
||||||
@@ -80,6 +272,7 @@ def run_freshrss_pipeline(
|
|||||||
debug_artifacts: bool = False,
|
debug_artifacts: bool = False,
|
||||||
prompt: Path | None = None,
|
prompt: Path | None = None,
|
||||||
rules: Path | None = None,
|
rules: Path | None = None,
|
||||||
|
context: dict[str, Any] | None = None,
|
||||||
context_path: Path | None = None,
|
context_path: Path | None = None,
|
||||||
max_retries: int = 2,
|
max_retries: int = 2,
|
||||||
timeout_seconds: float = 60.0,
|
timeout_seconds: float = 60.0,
|
||||||
@@ -90,6 +283,7 @@ def run_freshrss_pipeline(
|
|||||||
delivery_date: date | None = None,
|
delivery_date: date | None = None,
|
||||||
output_dir: Path | None = None,
|
output_dir: Path | None = None,
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
|
# --- 阶段 1:初始化 run_id、输出路径、凭证 ---
|
||||||
started_at = datetime.now(tz=UTC)
|
started_at = datetime.now(tz=UTC)
|
||||||
resolved_output_dir = output_dir or default_output_dir()
|
resolved_output_dir = output_dir or default_output_dir()
|
||||||
run_stamp = started_at.strftime("%Y%m%d-%H%M%S")
|
run_stamp = started_at.strftime("%Y%m%d-%H%M%S")
|
||||||
@@ -112,6 +306,7 @@ def run_freshrss_pipeline(
|
|||||||
report_output = resolved_output_dir / "run-report.json"
|
report_output = resolved_output_dir / "run-report.json"
|
||||||
items_list_output = _maybe_path(debug_artifacts, resolved_output_dir / "items" / "freshrss.items.json")
|
items_list_output = _maybe_path(debug_artifacts, resolved_output_dir / "items" / "freshrss.items.json")
|
||||||
|
|
||||||
|
# --- 阶段 2:登录 FreshRSS,拉取未读条目,写原始 payload ---
|
||||||
client = FreshRSSClient(
|
client = FreshRSSClient(
|
||||||
api_base_url=resolved_api_base_url,
|
api_base_url=resolved_api_base_url,
|
||||||
username=resolved_username,
|
username=resolved_username,
|
||||||
@@ -137,158 +332,71 @@ def run_freshrss_pipeline(
|
|||||||
_save_json(items_list_output, [item.model_dump(mode="json") for item in items])
|
_save_json(items_list_output, [item.model_dump(mode="json") for item in items])
|
||||||
|
|
||||||
loaded_rules = load_filter_rules(resolved_rules_path)
|
loaded_rules = load_filter_rules(resolved_rules_path)
|
||||||
context = FilterContext.model_validate(_load_json(context_path)) if context_path else FilterContext()
|
# context 优先使用直接传入的 dict,其次读取 context_path 文件,两者均缺失则使用空 context
|
||||||
|
if context is not None:
|
||||||
|
filter_context = FilterContext.model_validate(context)
|
||||||
|
elif context_path is not None:
|
||||||
|
filter_context = FilterContext.model_validate(_load_json(context_path))
|
||||||
|
else:
|
||||||
|
filter_context = FilterContext()
|
||||||
|
|
||||||
|
# --- 阶段 3:逐条处理(提取 -> LLM 摘要 -> 规则过滤 -> 候选构建) ---
|
||||||
delivered_candidates: list[OpenClawCandidateInput] = []
|
delivered_candidates: list[OpenClawCandidateInput] = []
|
||||||
delivered_item_ids: list[str] = []
|
delivered_item_ids: list[str] = []
|
||||||
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
|
result = _process_item(
|
||||||
item_key = f"item-{index:02d}"
|
index=index,
|
||||||
item_path = _maybe_path(debug_artifacts, resolved_output_dir / "items" / f"{item_key}.item.json")
|
item=item,
|
||||||
extracted_path = resolved_output_dir / "extracted" / f"{item_key}.extracted.json"
|
resolved_output_dir=resolved_output_dir,
|
||||||
summary_output = _maybe_path(debug_artifacts, resolved_output_dir / "summary" / item_key / "result.loop.json")
|
resolved_prompt_path=resolved_prompt_path,
|
||||||
filter_path = _maybe_path(debug_artifacts, resolved_output_dir / "filter" / f"{item_key}.filter.json")
|
resolved_run_id=resolved_run_id,
|
||||||
record_path = _maybe_path(debug_artifacts, resolved_output_dir / "candidates" / f"{item_key}.article-candidate-record.json")
|
debug_artifacts=debug_artifacts,
|
||||||
openclaw_path = _maybe_path(debug_artifacts, resolved_output_dir / "candidates" / f"{item_key}.openclaw-candidate-input.json")
|
loaded_rules=loaded_rules,
|
||||||
|
filter_context=filter_context,
|
||||||
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,
|
|
||||||
max_retries=max_retries,
|
max_retries=max_retries,
|
||||||
timeout_seconds=timeout_seconds,
|
timeout_seconds=timeout_seconds,
|
||||||
api_key=resolved_llm_api_key,
|
resolved_llm_api_key=resolved_llm_api_key,
|
||||||
model=resolved_llm_model,
|
resolved_llm_model=resolved_llm_model,
|
||||||
api_url=resolved_llm_api_url,
|
resolved_llm_api_url=resolved_llm_api_url,
|
||||||
)
|
)
|
||||||
if summary_exit_code != 0 or summary_payload is None:
|
if result.get("status") == "delivered":
|
||||||
item_report["status"] = "summary_failed"
|
delivered_candidates.append(result.pop("_candidate"))
|
||||||
if summary_report is not None:
|
external_id = result.pop("_external_id", None)
|
||||||
item_report["summary_errors"] = summary_report.errors
|
if external_id:
|
||||||
item_reports.append(item_report)
|
delivered_item_ids.append(external_id)
|
||||||
continue
|
item_reports.append(result)
|
||||||
|
|
||||||
summary = LlmSummaryResult.model_validate(summary_payload)
|
# --- 阶段 4+5:构建 delivery payload 并持久化词元索引 ---
|
||||||
decision = evaluate_filter_rules(
|
delivery_payload, keyword_index_result = _build_and_persist_delivery(
|
||||||
FilterInput(item=item, article=extraction.article, summary=summary, context=context),
|
delivered_candidates=delivered_candidates,
|
||||||
loaded_rules,
|
resolved_run_id=resolved_run_id,
|
||||||
)
|
resolved_delivery_date=resolved_delivery_date,
|
||||||
if filter_path is not None:
|
delivery_output=delivery_output,
|
||||||
_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,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# --- 阶段 6:标记已读,汇总报告,返回结果 ---
|
||||||
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 保持未读状态
|
# 仅标记成功投递(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})
|
||||||
|
|
||||||
status_counts: dict[str, int] = {}
|
report = _build_run_report(
|
||||||
for item_report in item_reports:
|
resolved_run_id=resolved_run_id,
|
||||||
status = str(item_report["status"])
|
started_at=started_at,
|
||||||
status_counts[status] = status_counts.get(status, 0) + 1
|
limit=limit,
|
||||||
|
items=items,
|
||||||
report = {
|
delivered_candidates=delivered_candidates,
|
||||||
"run_id": resolved_run_id,
|
marked_count=marked_count,
|
||||||
"started_at": started_at.isoformat(),
|
mark_read=mark_read,
|
||||||
"completed_at": datetime.now(tz=UTC).isoformat(),
|
debug_artifacts=debug_artifacts,
|
||||||
"requested_limit": limit,
|
raw_output=raw_output,
|
||||||
"pulled_count": len(items),
|
delivery_output=delivery_output,
|
||||||
"delivered_count": len(delivered_candidates),
|
keyword_index_result=keyword_index_result,
|
||||||
"marked_read_count": marked_count,
|
item_reports=item_reports,
|
||||||
"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,
|
|
||||||
}
|
|
||||||
_save_json(report_output, report)
|
_save_json(report_output, report)
|
||||||
|
|
||||||
return {
|
return {
|
||||||
@@ -301,7 +409,7 @@ def run_freshrss_pipeline(
|
|||||||
"pulled_count": len(items),
|
"pulled_count": len(items),
|
||||||
"delivered_count": len(delivered_candidates),
|
"delivered_count": len(delivered_candidates),
|
||||||
"marked_read_count": marked_count,
|
"marked_read_count": marked_count,
|
||||||
"status_counts": status_counts,
|
"status_counts": report["status_counts"],
|
||||||
"debug_artifacts": debug_artifacts,
|
"debug_artifacts": debug_artifacts,
|
||||||
"delivery_payload": delivery_payload.model_dump(mode="json"),
|
"delivery_payload": delivery_payload.model_dump(mode="json"),
|
||||||
"items": item_reports,
|
"items": item_reports,
|
||||||
|
|||||||
Reference in New Issue
Block a user