feat: add freshrss openclaw pipeline and clean repo
This commit is contained in:
@@ -6,14 +6,24 @@ from summary_mcp.core.errors import SummaryError
|
||||
from summary_mcp.models.summary_io import ExtractionInput
|
||||
|
||||
|
||||
MIN_INLINE_CONTENT_LENGTH = 500
|
||||
|
||||
|
||||
def _has_usable_content(value: str | None) -> bool:
|
||||
return bool(value and len(value.strip()) >= MIN_INLINE_CONTENT_LENGTH)
|
||||
|
||||
|
||||
def choose_inline_content(extraction_input: ExtractionInput) -> tuple[str | None, str]:
|
||||
if extraction_input.raw_html:
|
||||
return extraction_input.raw_html, "raw_html"
|
||||
|
||||
if extraction_input.item and extraction_input.item.raw_content and len(extraction_input.item.raw_content.strip()) >= 500:
|
||||
if extraction_input.item and _has_usable_content(extraction_input.item.raw_content):
|
||||
return extraction_input.item.raw_content, "item.raw_content"
|
||||
|
||||
if extraction_input.rss_content and len(extraction_input.rss_content.strip()) >= 500:
|
||||
if extraction_input.item and _has_usable_content(extraction_input.item.raw_summary):
|
||||
return extraction_input.item.raw_summary, "item.raw_summary"
|
||||
|
||||
if _has_usable_content(extraction_input.rss_content):
|
||||
return extraction_input.rss_content, "rss_content"
|
||||
|
||||
return None, "none"
|
||||
|
||||
@@ -9,6 +9,19 @@ from summary_mcp.core.quality_checker import assess_quality
|
||||
from summary_mcp.models.summary_io import DebugInfo, ExtractionInput, ExtractionOutput
|
||||
|
||||
|
||||
RSS_ONLY_UPSTREAMS = {"freshrss"}
|
||||
|
||||
|
||||
def _should_skip_fetch(extraction_input: ExtractionInput) -> bool:
|
||||
item = extraction_input.item
|
||||
if item is None:
|
||||
return False
|
||||
|
||||
metadata = item.metadata if isinstance(item.metadata, dict) else {}
|
||||
upstream = metadata.get("upstream")
|
||||
return isinstance(upstream, str) and upstream in RSS_ONLY_UPSTREAMS
|
||||
|
||||
|
||||
def extract_content(extraction_input: ExtractionInput) -> ExtractionOutput:
|
||||
try:
|
||||
normalized = normalize_input(extraction_input)
|
||||
@@ -16,11 +29,20 @@ def extract_content(extraction_input: ExtractionInput) -> ExtractionOutput:
|
||||
|
||||
inline_content, content_source = choose_inline_content(normalized)
|
||||
if inline_content is None:
|
||||
if _should_skip_fetch(normalized):
|
||||
raise SummaryError(
|
||||
code="RSS_CONTENT_MISSING",
|
||||
message="Skipping item because RSS content is unavailable.",
|
||||
retryable=False,
|
||||
stage="extract",
|
||||
details={"url": str(normalized.item.url), "upstream": normalized.item.metadata.get("upstream")},
|
||||
)
|
||||
inline_content = fetch_html(str(normalized.item.url))
|
||||
content_source = "fetched_html"
|
||||
|
||||
if normalized.item.title is None and (
|
||||
content_source in {"raw_html", "fetched_html"} or inline_content.lstrip().startswith("<")
|
||||
content_source in {"raw_html", "fetched_html", "item.raw_content", "item.raw_summary", "rss_content"}
|
||||
or inline_content.lstrip().startswith("<")
|
||||
):
|
||||
normalized.item.title = extract_title(inline_content)
|
||||
|
||||
|
||||
@@ -0,0 +1,273 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
import httpx
|
||||
|
||||
from summary_mcp.validators.llm_result import ValidationReport
|
||||
from summary_mcp.validators.llm_result import validate_llm_result as validate_llm_result_from_path
|
||||
from summary_mcp.validators.llm_result import validate_llm_result_payload
|
||||
|
||||
|
||||
JSON_BLOCK_RE = re.compile(r"```(?:json)?\s*(\{.*\})\s*```", re.DOTALL)
|
||||
DEFAULT_CHAT_COMPLETIONS_URL = "https://api.openai.com/v1/chat/completions"
|
||||
|
||||
|
||||
def load_text(path: Path) -> str:
|
||||
return path.read_text(encoding="utf-8")
|
||||
|
||||
|
||||
def load_json(path: Path) -> dict[str, Any]:
|
||||
return json.loads(load_text(path))
|
||||
|
||||
|
||||
def save_json(path: Path, payload: dict[str, 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 build_summary_input(extracted: dict[str, Any]) -> dict[str, Any]:
|
||||
article = extracted.get("article") or {}
|
||||
return {
|
||||
"article": {
|
||||
"title": article.get("title"),
|
||||
"url": article.get("url"),
|
||||
"plain_text": article.get("plain_text"),
|
||||
"quality_flags": article.get("quality_flags"),
|
||||
},
|
||||
"warnings": extracted.get("warnings", []),
|
||||
}
|
||||
|
||||
|
||||
def build_initial_prompt(prompt_template: str, extracted: dict[str, Any]) -> str:
|
||||
summary_input = build_summary_input(extracted)
|
||||
return (
|
||||
f"{prompt_template}\n\n"
|
||||
"Below is the structured extracted article input. Generate the final summary JSON from it.\n\n"
|
||||
f"{json.dumps(summary_input, ensure_ascii=False, indent=2)}"
|
||||
)
|
||||
|
||||
|
||||
def build_repair_prompt(
|
||||
errors: list[str],
|
||||
extracted: dict[str, Any],
|
||||
result_json: dict[str, Any],
|
||||
) -> str:
|
||||
summary_input = build_summary_input(extracted)
|
||||
return (
|
||||
"Please repair the following invalid summary JSON.\n\n"
|
||||
"Requirements:\n"
|
||||
"- Output valid JSON only\n"
|
||||
"- Keep fields that are already correct\n"
|
||||
"- Fix only the validator-reported errors\n"
|
||||
"- Do not add explanations\n\n"
|
||||
f"validator errors:\n{json.dumps(errors, ensure_ascii=False, indent=2)}\n\n"
|
||||
f"Extracted article input:\n{json.dumps(summary_input, ensure_ascii=False, indent=2)}\n\n"
|
||||
f"Current summary JSON:\n{json.dumps(result_json, ensure_ascii=False, indent=2)}\n"
|
||||
)
|
||||
|
||||
|
||||
def extract_json_text(raw_text: str) -> str:
|
||||
fenced = JSON_BLOCK_RE.search(raw_text)
|
||||
if fenced:
|
||||
return fenced.group(1)
|
||||
|
||||
stripped = raw_text.strip()
|
||||
start = stripped.find("{")
|
||||
end = stripped.rfind("}")
|
||||
if start == -1 or end == -1 or end <= start:
|
||||
raise ValueError("Model output does not contain a JSON object.")
|
||||
return stripped[start : end + 1]
|
||||
|
||||
|
||||
def normalize_chat_completions_url(api_url: str | None) -> str | None:
|
||||
if api_url is None:
|
||||
return None
|
||||
|
||||
normalized = api_url.strip().rstrip("/")
|
||||
if not normalized:
|
||||
return None
|
||||
if normalized.endswith("/chat/completions"):
|
||||
return normalized
|
||||
return f"{normalized}/chat/completions"
|
||||
|
||||
|
||||
def resolve_llm_settings(
|
||||
*,
|
||||
api_key: str | None = None,
|
||||
model: str | None = None,
|
||||
api_url: str | None = None,
|
||||
) -> tuple[str, str, str]:
|
||||
resolved_api_key = api_key or os.environ.get("LLM_API_KEY") or os.environ.get("OPENAI_API_KEY")
|
||||
if not resolved_api_key:
|
||||
raise RuntimeError("Missing LLM_API_KEY or OPENAI_API_KEY, or pass an API key.")
|
||||
|
||||
resolved_model = model or os.environ.get("LLM_MODEL") or os.environ.get("OPENAI_MODEL")
|
||||
if not resolved_model:
|
||||
raise RuntimeError("Missing LLM_MODEL or OPENAI_MODEL, or pass a model.")
|
||||
|
||||
resolved_api_url = normalize_chat_completions_url(
|
||||
api_url or os.environ.get("LLM_API_URL") or os.environ.get("OPENAI_API_URL") or DEFAULT_CHAT_COMPLETIONS_URL
|
||||
)
|
||||
if not resolved_api_url:
|
||||
raise RuntimeError("Missing LLM_API_URL, OPENAI_API_URL, or pass an API URL.")
|
||||
|
||||
return resolved_api_key, resolved_model, resolved_api_url
|
||||
|
||||
|
||||
def call_llm(
|
||||
prompt: str,
|
||||
timeout_seconds: float,
|
||||
api_key: str | None,
|
||||
model: str | None,
|
||||
api_url: str | None,
|
||||
) -> str:
|
||||
resolved_api_key, resolved_model, resolved_api_url = resolve_llm_settings(
|
||||
api_key=api_key,
|
||||
model=model,
|
||||
api_url=api_url,
|
||||
)
|
||||
|
||||
headers = {
|
||||
"Authorization": f"Bearer {resolved_api_key}",
|
||||
"Content-Type": "application/json",
|
||||
}
|
||||
payload = {
|
||||
"model": resolved_model,
|
||||
"messages": [
|
||||
{
|
||||
"role": "system",
|
||||
"content": "You are a precise JSON generator. Always output a single valid JSON object.",
|
||||
},
|
||||
{"role": "user", "content": prompt},
|
||||
],
|
||||
"temperature": 0.2,
|
||||
}
|
||||
|
||||
with httpx.Client(timeout=timeout_seconds) as client:
|
||||
response = client.post(resolved_api_url, headers=headers, json=payload)
|
||||
response.raise_for_status()
|
||||
data = response.json()
|
||||
|
||||
try:
|
||||
return data["choices"][0]["message"]["content"]
|
||||
except (KeyError, IndexError, TypeError) as exc:
|
||||
raise RuntimeError(f"Unexpected LLM response shape: {json.dumps(data, ensure_ascii=False)[:1000]}") from exc
|
||||
|
||||
|
||||
def _save_attempt_artifact(base_output_path: Path | None, suffix: str, payload: str | dict[str, Any]) -> None:
|
||||
if base_output_path is None:
|
||||
return
|
||||
|
||||
path = base_output_path.with_name(f"{base_output_path.stem}.{suffix}")
|
||||
if isinstance(payload, str):
|
||||
path.write_text(payload, encoding="utf-8")
|
||||
else:
|
||||
save_json(path, payload)
|
||||
|
||||
|
||||
def run_loop_payload(
|
||||
*,
|
||||
extracted_payload: dict[str, Any],
|
||||
prompt_path: Path,
|
||||
max_retries: int,
|
||||
timeout_seconds: float,
|
||||
api_key: str | None,
|
||||
model: str | None,
|
||||
api_url: str | None,
|
||||
output_path: Path | None = None,
|
||||
) -> tuple[int, dict[str, Any] | None, ValidationReport | None]:
|
||||
prompt_template = load_text(prompt_path)
|
||||
if output_path is not None:
|
||||
output_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
last_errors: list[str] = []
|
||||
last_result: dict[str, Any] | None = None
|
||||
|
||||
for attempt in range(1, max_retries + 2):
|
||||
if attempt == 1:
|
||||
prompt = build_initial_prompt(prompt_template, extracted_payload)
|
||||
else:
|
||||
assert last_result is not None
|
||||
prompt = build_repair_prompt(last_errors, extracted_payload, last_result)
|
||||
|
||||
raw_output = call_llm(prompt, timeout_seconds, api_key, model, api_url)
|
||||
_save_attempt_artifact(output_path, f"attempt-{attempt}.raw.txt", raw_output)
|
||||
|
||||
try:
|
||||
result_payload = json.loads(extract_json_text(raw_output))
|
||||
except (json.JSONDecodeError, ValueError) as exc:
|
||||
last_errors = [f"Model output is not valid JSON: {exc}"]
|
||||
last_result = {"raw_output": raw_output}
|
||||
report_payload = {
|
||||
"valid": False,
|
||||
"errors": last_errors,
|
||||
"warnings": [],
|
||||
"normalized_result": None,
|
||||
}
|
||||
_save_attempt_artifact(output_path, f"attempt-{attempt}.validation.json", report_payload)
|
||||
if attempt > max_retries:
|
||||
if output_path is not None:
|
||||
output_path.write_text(raw_output, encoding="utf-8")
|
||||
return 1, None, ValidationReport(valid=False, errors=last_errors)
|
||||
continue
|
||||
|
||||
_save_attempt_artifact(output_path, f"attempt-{attempt}.json", result_payload)
|
||||
if output_path is not None:
|
||||
save_json(output_path, result_payload)
|
||||
|
||||
report = validate_llm_result_payload(result_payload, extracted_payload)
|
||||
_save_attempt_artifact(
|
||||
output_path,
|
||||
f"attempt-{attempt}.validation.json",
|
||||
{
|
||||
"valid": report.valid,
|
||||
"errors": report.errors,
|
||||
"warnings": report.warnings,
|
||||
"normalized_result": report.normalized_result,
|
||||
},
|
||||
)
|
||||
|
||||
if report.valid:
|
||||
return 0, result_payload, report
|
||||
|
||||
last_errors = report.errors
|
||||
last_result = result_payload
|
||||
|
||||
return 1, None, ValidationReport(valid=False, errors=last_errors)
|
||||
|
||||
|
||||
def run_loop(
|
||||
extracted_path: Path,
|
||||
prompt_path: Path,
|
||||
output_path: Path,
|
||||
max_retries: int,
|
||||
timeout_seconds: float,
|
||||
api_key: str | None,
|
||||
model: str | None,
|
||||
api_url: str | None,
|
||||
) -> int:
|
||||
extracted = load_json(extracted_path)
|
||||
exit_code, _, report = run_loop_payload(
|
||||
extracted_payload=extracted,
|
||||
prompt_path=prompt_path,
|
||||
output_path=output_path,
|
||||
max_retries=max_retries,
|
||||
timeout_seconds=timeout_seconds,
|
||||
api_key=api_key,
|
||||
model=model,
|
||||
api_url=api_url,
|
||||
)
|
||||
|
||||
if exit_code == 0:
|
||||
return 0
|
||||
|
||||
if report is not None and output_path.exists() and extracted_path.exists():
|
||||
fallback_report = validate_llm_result_from_path(output_path, extracted_path)
|
||||
if fallback_report.valid:
|
||||
return 0
|
||||
return exit_code
|
||||
@@ -1,15 +1,20 @@
|
||||
from __future__ import annotations
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
from datetime import UTC, datetime
|
||||
from typing import Any
|
||||
from urllib.parse import urlparse
|
||||
from urllib.parse import urlencode, urlparse
|
||||
|
||||
import httpx
|
||||
|
||||
from summary_mcp.models.item import Item
|
||||
|
||||
|
||||
READ_TAG = "user/-/state/com.google/read"
|
||||
KEPT_UNREAD_TAG = "user/-/state/com.google/kept-unread"
|
||||
MIN_INLINE_CONTENT_LENGTH = 500
|
||||
|
||||
|
||||
def _trim_api_base_url(api_base_url: str) -> str:
|
||||
return api_base_url.rstrip("/")
|
||||
|
||||
@@ -100,6 +105,24 @@ def _parse_datetime(timestamp: Any) -> datetime | None:
|
||||
return None
|
||||
|
||||
|
||||
def _dedupe_non_empty(values: list[str] | None) -> list[str]:
|
||||
if not values:
|
||||
return []
|
||||
|
||||
seen: set[str] = set()
|
||||
deduped: list[str] = []
|
||||
for value in values:
|
||||
if not value or value in seen:
|
||||
continue
|
||||
seen.add(value)
|
||||
deduped.append(value)
|
||||
return deduped
|
||||
|
||||
|
||||
def _has_usable_inline_content(value: str | None) -> bool:
|
||||
return bool(value and len(value.strip()) >= MIN_INLINE_CONTENT_LENGTH)
|
||||
|
||||
|
||||
def map_entry_to_item(entry: dict[str, Any]) -> Item:
|
||||
url = _pick_entry_url(entry)
|
||||
summary_html = _pick_content_block(entry, "summary")
|
||||
@@ -107,7 +130,7 @@ def map_entry_to_item(entry: dict[str, Any]) -> Item:
|
||||
published_at = _parse_datetime(entry.get("published"))
|
||||
source_id = _build_source_id(entry, url)
|
||||
external_id = entry.get("id") if isinstance(entry.get("id"), str) else None
|
||||
fetch_state = "fetched" if content_html and len(content_html.strip()) >= 500 else "pending"
|
||||
fetch_state = "fetched" if _has_usable_inline_content(content_html) or _has_usable_inline_content(summary_html) else "pending"
|
||||
|
||||
metadata = {
|
||||
"upstream": "freshrss",
|
||||
@@ -151,6 +174,9 @@ class FreshRSSClient:
|
||||
def _build_url(self, path: str) -> str:
|
||||
return f"{self.api_base_url}/{path.lstrip('/')}"
|
||||
|
||||
def _build_auth_headers(self, auth_token: str) -> dict[str, str]:
|
||||
return {"Authorization": f"GoogleLogin auth={auth_token}"}
|
||||
|
||||
def client_login(self) -> str:
|
||||
data = {
|
||||
"Email": self.username,
|
||||
@@ -171,23 +197,42 @@ class FreshRSSClient:
|
||||
|
||||
return auth_token
|
||||
|
||||
def get_edit_token(self, auth_token: str) -> str:
|
||||
with httpx.Client(
|
||||
timeout=self.timeout_seconds,
|
||||
headers=self._build_auth_headers(auth_token),
|
||||
) as client:
|
||||
response = client.get(self._build_url("reader/api/0/token"))
|
||||
response.raise_for_status()
|
||||
|
||||
edit_token = response.text.strip()
|
||||
if not edit_token:
|
||||
raise RuntimeError("FreshRSS token request succeeded but did not return an edit token.")
|
||||
return edit_token
|
||||
|
||||
def fetch_stream_contents(
|
||||
self,
|
||||
auth_token: str,
|
||||
stream_id: str = "user/-/state/com.google/reading-list",
|
||||
limit: int = 20,
|
||||
continuation: str | None = None,
|
||||
exclude_targets: list[str] | None = None,
|
||||
include_targets: list[str] | None = None,
|
||||
) -> dict[str, Any]:
|
||||
params: dict[str, Any] = {
|
||||
"output": "json",
|
||||
"n": limit,
|
||||
}
|
||||
params: list[tuple[str, str | int]] = [
|
||||
("output", "json"),
|
||||
("n", limit),
|
||||
]
|
||||
if continuation:
|
||||
params["c"] = continuation
|
||||
params.append(("c", continuation))
|
||||
for target in _dedupe_non_empty(exclude_targets):
|
||||
params.append(("xt", target))
|
||||
for target in _dedupe_non_empty(include_targets):
|
||||
params.append(("it", target))
|
||||
|
||||
with httpx.Client(
|
||||
timeout=self.timeout_seconds,
|
||||
headers={"Authorization": f"GoogleLogin auth={auth_token}"},
|
||||
headers=self._build_auth_headers(auth_token),
|
||||
) as client:
|
||||
response = client.get(
|
||||
self._build_url(f"reader/api/0/stream/contents/{stream_id}"),
|
||||
@@ -196,4 +241,52 @@ class FreshRSSClient:
|
||||
response.raise_for_status()
|
||||
return response.json()
|
||||
|
||||
def edit_tag(
|
||||
self,
|
||||
auth_token: str,
|
||||
edit_token: str,
|
||||
item_ids: list[str],
|
||||
add_tags: list[str] | None = None,
|
||||
remove_tags: list[str] | None = None,
|
||||
) -> str:
|
||||
normalized_item_ids = _dedupe_non_empty(item_ids)
|
||||
if not normalized_item_ids:
|
||||
raise ValueError("FreshRSS edit-tag requires at least one item id.")
|
||||
|
||||
payload: list[tuple[str, str]] = [("ac", "edit"), ("T", edit_token)]
|
||||
for item_id in normalized_item_ids:
|
||||
payload.append(("i", item_id))
|
||||
for tag in _dedupe_non_empty(add_tags):
|
||||
payload.append(("a", tag))
|
||||
for tag in _dedupe_non_empty(remove_tags):
|
||||
payload.append(("r", tag))
|
||||
|
||||
headers = self._build_auth_headers(auth_token)
|
||||
headers["Content-Type"] = "application/x-www-form-urlencoded"
|
||||
|
||||
with httpx.Client(
|
||||
timeout=self.timeout_seconds,
|
||||
headers=headers,
|
||||
) as client:
|
||||
response = client.post(
|
||||
self._build_url("reader/api/0/edit-tag"),
|
||||
content=urlencode(payload).encode("utf-8"),
|
||||
)
|
||||
response.raise_for_status()
|
||||
|
||||
return response.text.strip()
|
||||
|
||||
def mark_items_as_read(
|
||||
self,
|
||||
auth_token: str,
|
||||
item_ids: list[str],
|
||||
edit_token: str | None = None,
|
||||
) -> str:
|
||||
resolved_edit_token = edit_token or self.get_edit_token(auth_token)
|
||||
return self.edit_tag(
|
||||
auth_token=auth_token,
|
||||
edit_token=resolved_edit_token,
|
||||
item_ids=item_ids,
|
||||
add_tags=[READ_TAG],
|
||||
remove_tags=[KEPT_UNREAD_TAG],
|
||||
)
|
||||
|
||||
@@ -1,4 +1,9 @@
|
||||
from __future__ import annotations
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import tempfile
|
||||
from datetime import date
|
||||
from pathlib import Path
|
||||
|
||||
from mcp.server.fastmcp import FastMCP
|
||||
|
||||
@@ -9,6 +14,7 @@ from summary_mcp.models.filtering import FilterContext, FilterInput
|
||||
from summary_mcp.models.item import Item
|
||||
from summary_mcp.models.llm_result import LlmSummaryResult
|
||||
from summary_mcp.models.summary_io import ExtractionInput
|
||||
from summary_mcp.workflows import run_freshrss_pipeline
|
||||
|
||||
|
||||
mcp = FastMCP(name="content-extract-mcp")
|
||||
@@ -30,9 +36,7 @@ def extract_url_content(url: str, language_hint: str | None = None) -> dict:
|
||||
def extract_item_content(item: dict) -> dict:
|
||||
"""Extract structured article content from a normalized item object."""
|
||||
parsed_item = Item.model_validate(item)
|
||||
result = extract_content(
|
||||
ExtractionInput(item=parsed_item)
|
||||
)
|
||||
result = extract_content(ExtractionInput(item=parsed_item))
|
||||
return result.model_dump(mode="json")
|
||||
|
||||
|
||||
@@ -61,6 +65,66 @@ def filter_summary_result(
|
||||
return decision.model_dump(mode="json")
|
||||
|
||||
|
||||
@mcp.tool()
|
||||
def run_freshrss_openclaw_pipeline(
|
||||
limit: int = 5,
|
||||
mark_read: bool = False,
|
||||
include_read: bool = False,
|
||||
debug_artifacts: bool = False,
|
||||
continuation: str | None = None,
|
||||
timeout_seconds: float = 60.0,
|
||||
max_retries: int = 2,
|
||||
stream_id: str = "user/-/state/com.google/reading-list",
|
||||
api_base_url: str | None = None,
|
||||
username: str | None = None,
|
||||
api_password: str | None = None,
|
||||
llm_api_key: str | None = None,
|
||||
llm_model: str | None = None,
|
||||
llm_api_url: str | None = None,
|
||||
context: dict | None = None,
|
||||
run_id: str | None = None,
|
||||
date_value: str | None = None,
|
||||
output_dir: str | None = None,
|
||||
include_item_reports: bool = False,
|
||||
) -> dict:
|
||||
"""Run the full FreshRSS -> extract -> LLM -> filter -> OpenClaw payload pipeline."""
|
||||
temp_context_path: Path | None = None
|
||||
|
||||
try:
|
||||
if context is not None:
|
||||
with tempfile.NamedTemporaryFile("w", encoding="utf-8", suffix=".json", delete=False) as handle:
|
||||
json.dump(context, handle, ensure_ascii=False, indent=2)
|
||||
temp_context_path = Path(handle.name)
|
||||
|
||||
result = run_freshrss_pipeline(
|
||||
api_base_url=api_base_url,
|
||||
username=username,
|
||||
api_password=api_password,
|
||||
stream_id=stream_id,
|
||||
limit=limit,
|
||||
continuation=continuation,
|
||||
include_read=include_read,
|
||||
mark_read=mark_read,
|
||||
debug_artifacts=debug_artifacts,
|
||||
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:
|
||||
result = {key: value for key, value in result.items() if key != "items"}
|
||||
return result
|
||||
|
||||
|
||||
def main() -> None:
|
||||
mcp.run()
|
||||
|
||||
|
||||
@@ -47,12 +47,31 @@ def _validate_business_rules(result: LlmSummaryResult, extracted: dict[str, Any]
|
||||
if extracted_url and str(result.url) != extracted_url:
|
||||
errors.append("`url` does not match extracted article url")
|
||||
|
||||
if result.category == "\u8d44\u8baf" and result.worth_keeping:
|
||||
warnings.append("`??` category marked as worth keeping; check if this is intentional.")
|
||||
if result.category == "资讯" and result.worth_keeping:
|
||||
warnings.append("`资讯` category marked as worth keeping; check if this is intentional.")
|
||||
|
||||
return errors, warnings
|
||||
|
||||
|
||||
def validate_llm_result_payload(
|
||||
result_payload: dict[str, Any],
|
||||
extracted_payload: dict[str, Any] | None = None,
|
||||
) -> ValidationReport:
|
||||
try:
|
||||
parsed = LlmSummaryResult.model_validate(result_payload)
|
||||
except ValidationError as exc:
|
||||
errors = [f"{'.'.join(str(part) for part in error['loc'])}: {error['msg']}" for error in exc.errors()]
|
||||
return ValidationReport(valid=False, errors=errors)
|
||||
|
||||
errors, warnings = _validate_business_rules(parsed, extracted_payload)
|
||||
return ValidationReport(
|
||||
valid=not errors,
|
||||
errors=errors,
|
||||
warnings=warnings,
|
||||
normalized_result=parsed.model_dump(mode="json"),
|
||||
)
|
||||
|
||||
|
||||
def validate_llm_result(
|
||||
result_path: Path,
|
||||
extracted_path: Path | None = None,
|
||||
@@ -69,16 +88,4 @@ def validate_llm_result(
|
||||
except json.JSONDecodeError as exc:
|
||||
return ValidationReport(valid=False, errors=[f"Invalid extracted JSON: {exc}"])
|
||||
|
||||
try:
|
||||
parsed = LlmSummaryResult.model_validate(raw_result)
|
||||
except ValidationError as exc:
|
||||
errors = [f"{'.'.join(str(part) for part in error['loc'])}: {error['msg']}" for error in exc.errors()]
|
||||
return ValidationReport(valid=False, errors=errors)
|
||||
|
||||
errors, warnings = _validate_business_rules(parsed, extracted)
|
||||
return ValidationReport(
|
||||
valid=not errors,
|
||||
errors=errors,
|
||||
warnings=warnings,
|
||||
normalized_result=parsed.model_dump(mode="json"),
|
||||
)
|
||||
return validate_llm_result_payload(raw_result, extracted)
|
||||
|
||||
@@ -0,0 +1,3 @@
|
||||
from .freshrss_pipeline import default_output_dir, read_delivery_payload, run_freshrss_pipeline
|
||||
|
||||
__all__ = ["default_output_dir", "read_delivery_payload", "run_freshrss_pipeline"]
|
||||
@@ -0,0 +1,284 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
from datetime import UTC, date, datetime
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
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"
|
||||
DEFAULT_PROMPT_PATH = OUTPUT_ROOT / "prompts" / "llm-summary-prompt.txt"
|
||||
DEFAULT_RULES_PATH = REPO_ROOT / "configs" / "filter_rules.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:
|
||||
resolved = value or os.environ.get(name)
|
||||
if not resolved:
|
||||
raise RuntimeError(f"Missing required value: {name}.")
|
||||
return resolved
|
||||
|
||||
|
||||
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 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_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]:
|
||||
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")
|
||||
|
||||
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 = FilterContext.model_validate(_load_json(context_path)) if context_path else FilterContext()
|
||||
|
||||
delivered_candidates: list[OpenClawCandidateInput] = []
|
||||
delivered_item_ids: list[str] = []
|
||||
item_reports: list[dict[str, Any]] = []
|
||||
|
||||
for index, item in enumerate(items, start=1):
|
||||
item_key = f"item-{index:02d}"
|
||||
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")
|
||||
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,
|
||||
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
|
||||
item_reports.append(item_report)
|
||||
continue
|
||||
|
||||
summary = LlmSummaryResult.model_validate(summary_payload)
|
||||
decision = evaluate_filter_rules(
|
||||
FilterInput(item=item, article=extraction.article, summary=summary, context=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"))
|
||||
|
||||
marked_count = 0
|
||||
if mark_read and 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})
|
||||
|
||||
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),
|
||||
"status_counts": status_counts,
|
||||
"items": 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),
|
||||
"pulled_count": len(items),
|
||||
"delivered_count": len(delivered_candidates),
|
||||
"marked_read_count": marked_count,
|
||||
"status_counts": 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))
|
||||
Reference in New Issue
Block a user