Refactor context passing: pipeline now accepts dict directly

Eliminate the NamedTemporaryFile workaround in server.py by adding a
`context` dict parameter to run_freshrss_pipeline. The pipeline resolves
filter context from the dict first, falling back to context_path, then
defaulting to an empty FilterContext. Remove unused json/tempfile imports.
This commit is contained in:
wdm
2026-03-29 16:27:11 +08:00
parent ced6142b59
commit 2d9a38e465
3 changed files with 30 additions and 37 deletions
+1 -1
View File
@@ -2,7 +2,7 @@
## 立即可改(低成本) ## 立即可改(低成本)
- [ ] `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 - [ ] `evaluate_filter_rules` 决策逻辑歧义 — `stop_on_match=True` 命中后应直接以该规则 decision 为最终结果,而非继续聚合所有 matched
+20 -34
View File
@@ -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"}
@@ -80,6 +80,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,
@@ -137,7 +138,13 @@ 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()
delivered_candidates: list[OpenClawCandidateInput] = [] delivered_candidates: list[OpenClawCandidateInput] = []
delivered_item_ids: list[str] = [] delivered_item_ids: list[str] = []
@@ -204,7 +211,7 @@ def run_freshrss_pipeline(
summary = LlmSummaryResult.model_validate(summary_payload) summary = LlmSummaryResult.model_validate(summary_payload)
decision = evaluate_filter_rules( decision = evaluate_filter_rules(
FilterInput(item=item, article=extraction.article, summary=summary, context=context), FilterInput(item=item, article=extraction.article, summary=summary, context=filter_context),
loaded_rules, loaded_rules,
) )
if filter_path is not None: if filter_path is not None: