Add FreshRSS ingestion and rule filtering
This commit is contained in:
@@ -0,0 +1,3 @@
|
||||
from .engine import DEFAULT_RULES_PATH, evaluate_filter_rules, load_filter_rules
|
||||
|
||||
__all__ = ["DEFAULT_RULES_PATH", "evaluate_filter_rules", "load_filter_rules"]
|
||||
@@ -0,0 +1,136 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from summary_mcp.models.filtering import (
|
||||
FieldCondition,
|
||||
FilterDecisionResult,
|
||||
FilterInput,
|
||||
FilterRule,
|
||||
MatchedRule,
|
||||
)
|
||||
|
||||
|
||||
REPO_ROOT = Path(__file__).resolve().parents[3]
|
||||
DEFAULT_RULES_PATH = REPO_ROOT / "configs" / "filter_rules.json"
|
||||
|
||||
|
||||
def load_filter_rules(path: Path | None = None) -> list[FilterRule]:
|
||||
rules_path = path or DEFAULT_RULES_PATH
|
||||
payload = json.loads(rules_path.read_text(encoding="utf-8-sig"))
|
||||
if not isinstance(payload, list):
|
||||
raise RuntimeError("Filter rules file must contain a JSON array.")
|
||||
return [FilterRule.model_validate(item) for item in payload]
|
||||
|
||||
|
||||
def _normalize_value(value: Any) -> Any:
|
||||
if hasattr(value, "model_dump"):
|
||||
return value.model_dump(mode="json")
|
||||
return value
|
||||
|
||||
|
||||
def _resolve_field(filter_input: FilterInput, field_path: str) -> Any:
|
||||
current: Any = filter_input
|
||||
for part in field_path.split("."):
|
||||
current = _normalize_value(current)
|
||||
if isinstance(current, dict):
|
||||
if part not in current:
|
||||
return None
|
||||
current = current[part]
|
||||
continue
|
||||
return None
|
||||
return _normalize_value(current)
|
||||
|
||||
|
||||
def _expected_value(filter_input: FilterInput, value: Any) -> Any:
|
||||
if isinstance(value, dict) and "from_field" in value:
|
||||
field_name = value.get("from_field")
|
||||
if isinstance(field_name, str):
|
||||
return _resolve_field(filter_input, field_name)
|
||||
return value
|
||||
|
||||
|
||||
def _match_condition(filter_input: FilterInput, condition: FieldCondition) -> bool:
|
||||
current = _resolve_field(filter_input, condition.field)
|
||||
expected = _expected_value(filter_input, condition.value)
|
||||
|
||||
if condition.op == "exists":
|
||||
return (current is not None) if expected is not False else (current is None)
|
||||
if condition.op == "eq":
|
||||
return current == expected
|
||||
if condition.op == "ne":
|
||||
return current != expected
|
||||
if condition.op == "in":
|
||||
return current in expected if isinstance(expected, list) else False
|
||||
if condition.op == "not_in":
|
||||
return current not in expected if isinstance(expected, list) else False
|
||||
if condition.op == "contains":
|
||||
if isinstance(current, list):
|
||||
return expected in current
|
||||
if isinstance(current, str) and isinstance(expected, str):
|
||||
return expected in current
|
||||
return False
|
||||
if condition.op == "overlap":
|
||||
if isinstance(current, list) and isinstance(expected, list):
|
||||
return bool(set(current) & set(expected))
|
||||
return False
|
||||
if condition.op == "gte":
|
||||
return current is not None and expected is not None and current >= expected
|
||||
if condition.op == "lte":
|
||||
return current is not None and expected is not None and current <= expected
|
||||
return False
|
||||
|
||||
|
||||
def _rule_matches(filter_input: FilterInput, rule: FilterRule) -> bool:
|
||||
if not rule.enabled:
|
||||
return False
|
||||
if rule.conditions_all and not all(_match_condition(filter_input, condition) for condition in rule.conditions_all):
|
||||
return False
|
||||
if rule.conditions_any and not any(_match_condition(filter_input, condition) for condition in rule.conditions_any):
|
||||
return False
|
||||
return bool(rule.conditions_all or rule.conditions_any)
|
||||
|
||||
|
||||
def evaluate_filter_rules(filter_input: FilterInput, rules: list[FilterRule]) -> FilterDecisionResult:
|
||||
matched: list[MatchedRule] = []
|
||||
|
||||
for rule in sorted(rules, key=lambda item: item.action.priority, reverse=True):
|
||||
if not _rule_matches(filter_input, rule):
|
||||
continue
|
||||
|
||||
matched_rule = MatchedRule(
|
||||
rule_id=rule.rule_id,
|
||||
decision=rule.action.decision,
|
||||
reason=rule.action.reason,
|
||||
labels=rule.action.labels,
|
||||
priority=rule.action.priority,
|
||||
)
|
||||
matched.append(matched_rule)
|
||||
if rule.stop_on_match:
|
||||
break
|
||||
|
||||
if any(rule.decision == "drop" for rule in matched):
|
||||
final_decision = "drop"
|
||||
elif any(rule.decision == "keep" for rule in matched):
|
||||
final_decision = "keep"
|
||||
elif any(rule.decision == "review" for rule in matched):
|
||||
final_decision = "review"
|
||||
else:
|
||||
final_decision = "review"
|
||||
|
||||
labels = sorted({label for rule in matched for label in rule.labels})
|
||||
reasons = [rule.reason for rule in matched]
|
||||
priorities = [rule.priority for rule in matched]
|
||||
|
||||
return FilterDecisionResult(
|
||||
decision=final_decision,
|
||||
matched_rules=[rule.rule_id for rule in matched],
|
||||
reasons=reasons if reasons else ["No rule matched; defaulted to review."],
|
||||
labels=labels,
|
||||
priority=max(priorities, default=0),
|
||||
matches=matched,
|
||||
)
|
||||
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
"""Integration helpers for upstream content sources."""
|
||||
@@ -0,0 +1,199 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
from datetime import UTC, datetime
|
||||
from typing import Any
|
||||
from urllib.parse import urlparse
|
||||
|
||||
import httpx
|
||||
|
||||
from summary_mcp.models.item import Item
|
||||
|
||||
|
||||
def _trim_api_base_url(api_base_url: str) -> str:
|
||||
return api_base_url.rstrip("/")
|
||||
|
||||
|
||||
def _build_item_id(source_id: str, external_id: str | None, url: str, published_at: datetime | None) -> str:
|
||||
seed = "|".join(
|
||||
[
|
||||
source_id,
|
||||
external_id or "",
|
||||
url,
|
||||
published_at.isoformat() if published_at else "",
|
||||
]
|
||||
)
|
||||
return f"sha256:{hashlib.sha256(seed.encode('utf-8')).hexdigest()}"
|
||||
|
||||
|
||||
def _build_source_id(entry: dict[str, Any], url: str) -> str:
|
||||
origin = entry.get("origin") or {}
|
||||
stream_id = origin.get("streamId")
|
||||
if isinstance(stream_id, str) and stream_id.strip():
|
||||
digest = hashlib.sha256(stream_id.encode("utf-8")).hexdigest()[:16]
|
||||
return f"freshrss:{digest}"
|
||||
|
||||
host = urlparse(url).netloc or "unknown-source"
|
||||
return f"freshrss:{host}"
|
||||
|
||||
|
||||
def _pick_entry_url(entry: dict[str, Any]) -> str:
|
||||
candidates = [
|
||||
entry.get("canonical"),
|
||||
entry.get("alternate"),
|
||||
]
|
||||
|
||||
for candidate_list in candidates:
|
||||
if not isinstance(candidate_list, list):
|
||||
continue
|
||||
for candidate in candidate_list:
|
||||
href = (candidate or {}).get("href")
|
||||
if isinstance(href, str) and href.strip():
|
||||
return href.strip()
|
||||
|
||||
entry_id = entry.get("id")
|
||||
if isinstance(entry_id, str) and entry_id.startswith("tag:"):
|
||||
return entry_id
|
||||
|
||||
raise ValueError("FreshRSS entry does not contain a usable URL.")
|
||||
|
||||
|
||||
def _pick_content_block(entry: dict[str, Any], key: str) -> str | None:
|
||||
block = entry.get(key)
|
||||
if isinstance(block, dict):
|
||||
content = block.get("content")
|
||||
if isinstance(content, str) and content.strip():
|
||||
return content
|
||||
return None
|
||||
|
||||
|
||||
def _pick_categories(entry: dict[str, Any]) -> list[str]:
|
||||
categories = entry.get("categories")
|
||||
if not isinstance(categories, list):
|
||||
return []
|
||||
|
||||
values: list[str] = []
|
||||
for category in categories:
|
||||
if not isinstance(category, str):
|
||||
continue
|
||||
if category.startswith("user/-/label/"):
|
||||
values.append(category.removeprefix("user/-/label/"))
|
||||
elif category.startswith("user/-/state/com.google/"):
|
||||
continue
|
||||
else:
|
||||
values.append(category)
|
||||
return values
|
||||
|
||||
|
||||
def _parse_datetime(timestamp: Any) -> datetime | None:
|
||||
if timestamp is None:
|
||||
return None
|
||||
if isinstance(timestamp, (int, float)):
|
||||
if timestamp > 10_000_000_000:
|
||||
return datetime.fromtimestamp(timestamp / 1000, tz=UTC)
|
||||
return datetime.fromtimestamp(timestamp, tz=UTC)
|
||||
if isinstance(timestamp, str) and timestamp.isdigit():
|
||||
value = int(timestamp)
|
||||
if value > 10_000_000_000:
|
||||
return datetime.fromtimestamp(value / 1000, tz=UTC)
|
||||
return datetime.fromtimestamp(value, tz=UTC)
|
||||
return None
|
||||
|
||||
|
||||
def map_entry_to_item(entry: dict[str, Any]) -> Item:
|
||||
url = _pick_entry_url(entry)
|
||||
summary_html = _pick_content_block(entry, "summary")
|
||||
content_html = _pick_content_block(entry, "content")
|
||||
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"
|
||||
|
||||
metadata = {
|
||||
"upstream": "freshrss",
|
||||
"origin": entry.get("origin") or {},
|
||||
"categories": _pick_categories(entry),
|
||||
"crawled_at": _parse_datetime(entry.get("crawlTimeMsec")),
|
||||
"published_epoch": entry.get("published"),
|
||||
}
|
||||
|
||||
return Item(
|
||||
item_id=_build_item_id(source_id, external_id, url, published_at),
|
||||
source_id=source_id,
|
||||
external_id=external_id,
|
||||
title=entry.get("title"),
|
||||
url=url,
|
||||
author=entry.get("author"),
|
||||
published_at=published_at,
|
||||
discovered_at=datetime.now(tz=UTC),
|
||||
content_kind="article",
|
||||
language=None,
|
||||
raw_summary=summary_html,
|
||||
raw_content=content_html,
|
||||
metadata=metadata,
|
||||
fetch_state=fetch_state,
|
||||
)
|
||||
|
||||
|
||||
class FreshRSSClient:
|
||||
def __init__(
|
||||
self,
|
||||
api_base_url: str,
|
||||
username: str,
|
||||
api_password: str,
|
||||
timeout_seconds: float = 20.0,
|
||||
) -> None:
|
||||
self.api_base_url = _trim_api_base_url(api_base_url)
|
||||
self.username = username
|
||||
self.api_password = api_password
|
||||
self.timeout_seconds = timeout_seconds
|
||||
|
||||
def _build_url(self, path: str) -> str:
|
||||
return f"{self.api_base_url}/{path.lstrip('/')}"
|
||||
|
||||
def client_login(self) -> str:
|
||||
data = {
|
||||
"Email": self.username,
|
||||
"Passwd": self.api_password,
|
||||
}
|
||||
with httpx.Client(timeout=self.timeout_seconds) as client:
|
||||
response = client.post(self._build_url("accounts/ClientLogin"), data=data)
|
||||
response.raise_for_status()
|
||||
|
||||
auth_token: str | None = None
|
||||
for line in response.text.splitlines():
|
||||
if line.startswith("Auth="):
|
||||
auth_token = line.split("=", 1)[1].strip()
|
||||
break
|
||||
|
||||
if not auth_token:
|
||||
raise RuntimeError("FreshRSS ClientLogin succeeded but did not return an Auth token.")
|
||||
|
||||
return auth_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,
|
||||
) -> dict[str, Any]:
|
||||
params: dict[str, Any] = {
|
||||
"output": "json",
|
||||
"n": limit,
|
||||
}
|
||||
if continuation:
|
||||
params["c"] = continuation
|
||||
|
||||
with httpx.Client(
|
||||
timeout=self.timeout_seconds,
|
||||
headers={"Authorization": f"GoogleLogin auth={auth_token}"},
|
||||
) as client:
|
||||
response = client.get(
|
||||
self._build_url(f"reader/api/0/stream/contents/{stream_id}"),
|
||||
params=params,
|
||||
)
|
||||
response.raise_for_status()
|
||||
return response.json()
|
||||
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
from typing import Any, Literal
|
||||
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
from .document import ExtractedArticle
|
||||
from .item import Item
|
||||
from .llm_result import LlmSummaryResult
|
||||
|
||||
|
||||
FilterDecision = Literal["keep", "drop", "review"]
|
||||
ConditionOp = Literal["eq", "ne", "in", "not_in", "contains", "overlap", "gte", "lte", "exists"]
|
||||
|
||||
|
||||
class FilterContext(BaseModel):
|
||||
source_tags: list[str] = Field(default_factory=list)
|
||||
interest_topics: list[str] = Field(default_factory=list)
|
||||
interest_keywords: list[str] = Field(default_factory=list)
|
||||
now: datetime | None = None
|
||||
|
||||
|
||||
class FieldCondition(BaseModel):
|
||||
field: str
|
||||
op: ConditionOp
|
||||
value: Any = None
|
||||
|
||||
|
||||
class FilterAction(BaseModel):
|
||||
decision: FilterDecision
|
||||
reason: str
|
||||
labels: list[str] = Field(default_factory=list)
|
||||
priority: int = 50
|
||||
|
||||
|
||||
class FilterRule(BaseModel):
|
||||
rule_id: str
|
||||
enabled: bool = True
|
||||
stop_on_match: bool = False
|
||||
conditions_all: list[FieldCondition] = Field(default_factory=list)
|
||||
conditions_any: list[FieldCondition] = Field(default_factory=list)
|
||||
action: FilterAction
|
||||
|
||||
|
||||
class FilterInput(BaseModel):
|
||||
item: Item | None = None
|
||||
article: ExtractedArticle | None = None
|
||||
summary: LlmSummaryResult
|
||||
context: FilterContext = Field(default_factory=FilterContext)
|
||||
|
||||
|
||||
class MatchedRule(BaseModel):
|
||||
rule_id: str
|
||||
decision: FilterDecision
|
||||
reason: str
|
||||
labels: list[str] = Field(default_factory=list)
|
||||
priority: int = 0
|
||||
|
||||
|
||||
class FilterDecisionResult(BaseModel):
|
||||
decision: FilterDecision
|
||||
matched_rules: list[str] = Field(default_factory=list)
|
||||
reasons: list[str] = Field(default_factory=list)
|
||||
labels: list[str] = Field(default_factory=list)
|
||||
priority: int = 0
|
||||
matches: list[MatchedRule] = Field(default_factory=list)
|
||||
@@ -1,9 +1,13 @@
|
||||
from __future__ import annotations
|
||||
from __future__ import annotations
|
||||
|
||||
from mcp.server.fastmcp import FastMCP
|
||||
|
||||
from summary_mcp.core.pipeline import extract_content
|
||||
from summary_mcp.filters.engine import evaluate_filter_rules, load_filter_rules
|
||||
from summary_mcp.models.document import ExtractedArticle
|
||||
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
|
||||
|
||||
|
||||
@@ -32,6 +36,31 @@ def extract_item_content(item: dict) -> dict:
|
||||
return result.model_dump(mode="json")
|
||||
|
||||
|
||||
@mcp.tool()
|
||||
def filter_summary_result(
|
||||
summary_result: dict,
|
||||
extracted_article: dict | None = None,
|
||||
item: dict | None = None,
|
||||
context: dict | None = None,
|
||||
) -> dict:
|
||||
"""Apply deterministic filter rules to a structured summary result."""
|
||||
parsed_summary = LlmSummaryResult.model_validate(summary_result)
|
||||
parsed_article = ExtractedArticle.model_validate(extracted_article) if extracted_article else None
|
||||
parsed_item = Item.model_validate(item) if item else None
|
||||
parsed_context = FilterContext.model_validate(context or {})
|
||||
rules = load_filter_rules()
|
||||
decision = evaluate_filter_rules(
|
||||
FilterInput(
|
||||
item=parsed_item,
|
||||
article=parsed_article,
|
||||
summary=parsed_summary,
|
||||
context=parsed_context,
|
||||
),
|
||||
rules,
|
||||
)
|
||||
return decision.model_dump(mode="json")
|
||||
|
||||
|
||||
def main() -> None:
|
||||
mcp.run()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user