Files
reader/src/summary_mcp/workflows/freshrss_pipeline.py
T

421 lines
17 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
from __future__ import annotations
# FreshRSS 全链路管道:拉取未读条目 -> 内容提取 -> LLM 摘要 -> 规则过滤 ->
# 构建 OpenClaw delivery payload -> 写盘 -> 词元统计 -> 标记已读。
# 生产入口:run_freshrss_pipeline(),由 MCP 工具 run_freshrss_openclaw_pipeline 调用。
import json
import os
from datetime import date, datetime, timezone
UTC = timezone.utc
from pathlib import Path
from typing import Any
from summary_mcp.core.keyword_index import persist_keyword_indexes
from summary_mcp.core.pipeline import extract_content
from summary_mcp.core.summary_loop import resolve_llm_settings, run_loop_payload
from summary_mcp.filters.engine import evaluate_filter_rules, load_filter_rules
from summary_mcp.integrations.freshrss import READ_TAG, FreshRSSClient, map_entry_to_item
from summary_mcp.models.article_candidate import (
CandidateMetadata,
CandidateSourceRefs,
OpenClawCandidateInput,
build_article_candidate_record,
build_openclaw_candidate_input,
)
from summary_mcp.models.filtering import FilterContext, FilterInput
from summary_mcp.models.llm_result import LlmSummaryResult
from summary_mcp.models.openclaw_delivery import OpenClawDeliveryPayload, build_openclaw_delivery_payload
from summary_mcp.models.summary_io import ExtractionInput
REPO_ROOT = Path(__file__).resolve().parents[3]
OUTPUT_ROOT = REPO_ROOT / "outputs"
FRESHRSS_OUTPUT_ROOT = OUTPUT_ROOT / "freshrss"
DATA_ROOT = REPO_ROOT / "data" / "term_index"
DEFAULT_PROMPT_PATH = OUTPUT_ROOT / "prompts" / "llm-summary-prompt.txt"
DEFAULT_RULES_PATH = REPO_ROOT / "configs" / "filter_rules.json"
DEFAULT_TERM_ALIASES_PATH = REPO_ROOT / "configs" / "term_aliases.json"
DEFAULT_TERM_STOPWORDS_PATH = REPO_ROOT / "configs" / "term_stopwords.json"
DEFAULT_TERM_DAILY_DIR = DATA_ROOT / "daily"
DEFAULT_TERM_STATS_PATH = DATA_ROOT / "term_stats.json"
def _save_json(path: Path, payload: dict[str, Any] | list[Any]) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
def _load_json(path: Path) -> dict[str, Any]:
return json.loads(path.read_text(encoding="utf-8-sig"))
def _load_required_env(name: str, value: str | None) -> str:
# 优先使用传入的 value,否则读取同名环境变量;两者均缺失时抛出 RuntimeError
if value:
return value
env_value = os.environ.get(name)
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:
return FRESHRSS_OUTPUT_ROOT / "rerun" / datetime.now(tz=UTC).strftime("%Y%m%d-%H%M%S")
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 供调用方解包
# 步骤 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(
*,
api_base_url: str | None = None,
username: str | None = None,
api_password: str | None = None,
stream_id: str = "user/-/state/com.google/reading-list",
limit: int = 5,
continuation: str | None = None,
include_read: bool = False,
mark_read: bool = False,
debug_artifacts: bool = False,
prompt: Path | None = None,
rules: Path | None = None,
context: dict[str, Any] | None = None,
context_path: Path | None = None,
max_retries: int = 2,
timeout_seconds: float = 60.0,
llm_api_key: str | None = None,
llm_model: str | None = None,
llm_api_url: str | None = None,
run_id: str | None = None,
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")
resolved_run_id = run_id or f"freshrss-pipeline-{run_stamp}"
resolved_delivery_date = delivery_date or datetime.now(tz=UTC).date()
resolved_api_base_url = _load_required_env("FRESHRSS_API_BASE_URL", api_base_url)
resolved_username = _load_required_env("FRESHRSS_USERNAME", username)
resolved_api_password = _load_required_env("FRESHRSS_API_PASSWORD", api_password)
resolved_llm_api_key, resolved_llm_model, resolved_llm_api_url = resolve_llm_settings(
api_key=llm_api_key,
model=llm_model,
api_url=llm_api_url,
)
resolved_prompt_path = prompt or DEFAULT_PROMPT_PATH
resolved_rules_path = rules or DEFAULT_RULES_PATH
raw_output = resolved_output_dir / "raw" / "freshrss.raw.json"
delivery_output = resolved_output_dir / "candidates" / "openclaw-delivery-payload.json"
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,
api_password=resolved_api_password,
timeout_seconds=timeout_seconds,
)
auth_token = client.client_login()
payload = client.fetch_stream_contents(
auth_token=auth_token,
stream_id=stream_id,
limit=limit,
continuation=continuation,
exclude_targets=[] if include_read else [READ_TAG],
)
entries = payload.get("items")
if not isinstance(entries, list):
raise RuntimeError("FreshRSS stream response does not contain an items array.")
_save_json(raw_output, payload)
items = [map_entry_to_item(entry) for entry in entries]
if items_list_output is not None:
_save_json(items_list_output, [item.model_dump(mode="json") for item in items])
loaded_rules = load_filter_rules(resolved_rules_path)
# 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_item_ids: list[str] = []
item_reports: list[dict[str, Any]] = []
for index, item in enumerate(items, start=1):
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,
resolved_llm_api_key=resolved_llm_api_key,
resolved_llm_model=resolved_llm_model,
resolved_llm_api_url=resolved_llm_api_url,
)
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)
# --- 阶段 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})
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 {
"run_id": resolved_run_id,
"output_dir": str(resolved_output_dir),
"raw_output": str(raw_output),
"delivery_output": str(delivery_output),
"report_output": str(report_output),
"keyword_index": keyword_index_result,
"pulled_count": len(items),
"delivered_count": len(delivered_candidates),
"marked_read_count": marked_count,
"status_counts": report["status_counts"],
"debug_artifacts": debug_artifacts,
"delivery_payload": delivery_payload.model_dump(mode="json"),
"items": item_reports,
}
def read_delivery_payload(path: Path) -> OpenClawDeliveryPayload:
return OpenClawDeliveryPayload.model_validate(_load_json(path))