Add keyword cleanup governance workflow
This commit is contained in:
@@ -0,0 +1,228 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from collections import Counter
|
||||
from datetime import UTC, date, datetime
|
||||
from pathlib import Path
|
||||
from typing import Iterable
|
||||
|
||||
from summary_mcp.models.article_candidate import OpenClawCandidateInput
|
||||
from summary_mcp.models.keyword_index import (
|
||||
DailyKeywordIndex,
|
||||
DailyKeywordTerm,
|
||||
KeywordStat,
|
||||
KeywordStatsIndex,
|
||||
)
|
||||
|
||||
|
||||
REPO_ROOT = Path(__file__).resolve().parents[3]
|
||||
DATA_ROOT = REPO_ROOT / "data" / "term_index"
|
||||
DEFAULT_DAILY_DIR = DATA_ROOT / "daily"
|
||||
DEFAULT_STATS_PATH = DATA_ROOT / "term_stats.json"
|
||||
DEFAULT_ALIASES_PATH = REPO_ROOT / "configs" / "term_aliases.json"
|
||||
DEFAULT_STOPWORDS_PATH = REPO_ROOT / "configs" / "term_stopwords.json"
|
||||
DEFAULT_INCLUDE_DECISIONS = {"keep", "review"}
|
||||
|
||||
|
||||
def _load_json(path: Path) -> object:
|
||||
return json.loads(path.read_text(encoding="utf-8-sig"))
|
||||
|
||||
|
||||
def _save_json(path: Path, payload: dict | list) -> None:
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
path.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
|
||||
|
||||
def _term_key(value: str) -> str:
|
||||
return value.strip().casefold()
|
||||
|
||||
|
||||
def load_term_aliases(path: Path | None = None) -> dict[str, str]:
|
||||
aliases_path = path or DEFAULT_ALIASES_PATH
|
||||
if not aliases_path.exists():
|
||||
return {}
|
||||
|
||||
payload = _load_json(aliases_path)
|
||||
if not isinstance(payload, dict):
|
||||
raise RuntimeError("Term aliases file must contain a JSON object.")
|
||||
|
||||
aliases: dict[str, str] = {}
|
||||
for raw_key, raw_value in payload.items():
|
||||
if not isinstance(raw_key, str) or not isinstance(raw_value, str):
|
||||
continue
|
||||
normalized_key = _term_key(raw_key)
|
||||
normalized_value = raw_value.strip()
|
||||
if not normalized_key or not normalized_value:
|
||||
continue
|
||||
aliases[normalized_key] = normalized_value
|
||||
return aliases
|
||||
|
||||
|
||||
def load_term_stopwords(path: Path | None = None) -> set[str]:
|
||||
stopwords_path = path or DEFAULT_STOPWORDS_PATH
|
||||
if not stopwords_path.exists():
|
||||
return set()
|
||||
|
||||
payload = _load_json(stopwords_path)
|
||||
if not isinstance(payload, list):
|
||||
raise RuntimeError("Term stopwords file must contain a JSON array.")
|
||||
|
||||
values: set[str] = set()
|
||||
for item in payload:
|
||||
if not isinstance(item, str):
|
||||
continue
|
||||
normalized = _term_key(item)
|
||||
if normalized:
|
||||
values.add(normalized)
|
||||
return values
|
||||
|
||||
|
||||
def normalize_keyword(
|
||||
keyword: str,
|
||||
*,
|
||||
aliases: dict[str, str],
|
||||
stopwords: set[str],
|
||||
) -> str | None:
|
||||
raw_value = keyword.strip()
|
||||
if not raw_value:
|
||||
return None
|
||||
|
||||
aliased_value = aliases.get(_term_key(raw_value), raw_value).strip()
|
||||
if not aliased_value:
|
||||
return None
|
||||
if _term_key(aliased_value) in stopwords:
|
||||
return None
|
||||
return aliased_value
|
||||
|
||||
|
||||
def build_daily_keyword_index(
|
||||
candidates: Iterable[OpenClawCandidateInput],
|
||||
*,
|
||||
for_date: date,
|
||||
digest_id: str,
|
||||
source: str = "openclaw_delivery_payload",
|
||||
include_decisions: set[str] | None = None,
|
||||
aliases: dict[str, str] | None = None,
|
||||
stopwords: set[str] | None = None,
|
||||
) -> DailyKeywordIndex:
|
||||
allowed_decisions = include_decisions or DEFAULT_INCLUDE_DECISIONS
|
||||
resolved_aliases = aliases or {}
|
||||
resolved_stopwords = stopwords or set()
|
||||
|
||||
term_counter: Counter[str] = Counter()
|
||||
candidate_count = 0
|
||||
|
||||
for candidate in candidates:
|
||||
if candidate.selection_decision not in allowed_decisions:
|
||||
continue
|
||||
|
||||
candidate_count += 1
|
||||
seen_for_candidate: set[str] = set()
|
||||
for keyword in candidate.keywords:
|
||||
normalized = normalize_keyword(
|
||||
keyword,
|
||||
aliases=resolved_aliases,
|
||||
stopwords=resolved_stopwords,
|
||||
)
|
||||
if normalized is None or normalized in seen_for_candidate:
|
||||
continue
|
||||
seen_for_candidate.add(normalized)
|
||||
term_counter[normalized] += 1
|
||||
|
||||
terms = [
|
||||
DailyKeywordTerm(term=term, normalized_term=term, count=count)
|
||||
for term, count in sorted(term_counter.items(), key=lambda item: (-item[1], item[0].casefold(), item[0]))
|
||||
]
|
||||
|
||||
return DailyKeywordIndex(
|
||||
schema_version="v1",
|
||||
date=for_date,
|
||||
source=source,
|
||||
digest_id=digest_id,
|
||||
generated_at=datetime.now(tz=UTC),
|
||||
candidate_count=candidate_count,
|
||||
terms=terms,
|
||||
)
|
||||
|
||||
|
||||
def read_daily_keyword_index(path: Path) -> DailyKeywordIndex:
|
||||
return DailyKeywordIndex.model_validate(_load_json(path))
|
||||
|
||||
|
||||
def rebuild_keyword_stats(*, daily_dir: Path = DEFAULT_DAILY_DIR) -> KeywordStatsIndex:
|
||||
term_totals: dict[str, int] = {}
|
||||
first_seen: dict[str, date] = {}
|
||||
last_seen: dict[str, date] = {}
|
||||
days_seen: dict[str, int] = {}
|
||||
|
||||
if daily_dir.exists():
|
||||
for path in sorted(daily_dir.glob("*.json")):
|
||||
daily_index = read_daily_keyword_index(path)
|
||||
seen_today: set[str] = set()
|
||||
for term in daily_index.terms:
|
||||
normalized = term.normalized_term
|
||||
term_totals[normalized] = term_totals.get(normalized, 0) + term.count
|
||||
if normalized not in first_seen or daily_index.date < first_seen[normalized]:
|
||||
first_seen[normalized] = daily_index.date
|
||||
if normalized not in last_seen or daily_index.date > last_seen[normalized]:
|
||||
last_seen[normalized] = daily_index.date
|
||||
if normalized not in seen_today:
|
||||
days_seen[normalized] = days_seen.get(normalized, 0) + 1
|
||||
seen_today.add(normalized)
|
||||
|
||||
terms = [
|
||||
KeywordStat(
|
||||
term=term,
|
||||
total_count=term_totals[term],
|
||||
days_seen=days_seen[term],
|
||||
first_seen=first_seen[term],
|
||||
last_seen=last_seen[term],
|
||||
)
|
||||
for term in sorted(term_totals, key=lambda item: (-term_totals[item], item.casefold(), item))
|
||||
]
|
||||
|
||||
return KeywordStatsIndex(
|
||||
schema_version="v1",
|
||||
generated_at=datetime.now(tz=UTC),
|
||||
terms=terms,
|
||||
)
|
||||
|
||||
|
||||
def persist_keyword_indexes(
|
||||
candidates: Iterable[OpenClawCandidateInput],
|
||||
*,
|
||||
for_date: date,
|
||||
digest_id: str,
|
||||
source: str = "openclaw_delivery_payload",
|
||||
daily_dir: Path = DEFAULT_DAILY_DIR,
|
||||
stats_path: Path = DEFAULT_STATS_PATH,
|
||||
aliases_path: Path | None = None,
|
||||
stopwords_path: Path | None = None,
|
||||
include_decisions: set[str] | None = None,
|
||||
) -> dict[str, object]:
|
||||
aliases = load_term_aliases(aliases_path)
|
||||
stopwords = load_term_stopwords(stopwords_path)
|
||||
daily_index = build_daily_keyword_index(
|
||||
candidates,
|
||||
for_date=for_date,
|
||||
digest_id=digest_id,
|
||||
source=source,
|
||||
include_decisions=include_decisions,
|
||||
aliases=aliases,
|
||||
stopwords=stopwords,
|
||||
)
|
||||
|
||||
daily_path = daily_dir / f"{for_date.isoformat()}.json"
|
||||
_save_json(daily_path, daily_index.model_dump(mode="json"))
|
||||
|
||||
stats_index = rebuild_keyword_stats(daily_dir=daily_dir)
|
||||
_save_json(stats_path, stats_index.model_dump(mode="json"))
|
||||
|
||||
return {
|
||||
"daily_output": str(daily_path),
|
||||
"stats_output": str(stats_path),
|
||||
"source": source,
|
||||
"candidate_count": daily_index.candidate_count,
|
||||
"term_count": len(daily_index.terms),
|
||||
"top_terms": [term.model_dump(mode="json") for term in daily_index.terms[:10]],
|
||||
}
|
||||
@@ -22,6 +22,7 @@ from .daily_digest import (
|
||||
DailyDigestSourceRef,
|
||||
DailyDigestStats,
|
||||
)
|
||||
from .keyword_index import DailyKeywordIndex, DailyKeywordTerm, KeywordStat, KeywordStatsIndex
|
||||
from .openclaw_delivery import (
|
||||
OpenClawDeliveryPayload,
|
||||
OpenClawDeliveryStats,
|
||||
@@ -40,7 +41,11 @@ __all__ = [
|
||||
"DailyDigestSection",
|
||||
"DailyDigestSourceRef",
|
||||
"DailyDigestStats",
|
||||
"DailyKeywordIndex",
|
||||
"DailyKeywordTerm",
|
||||
"DigestSectionHint",
|
||||
"KeywordStat",
|
||||
"KeywordStatsIndex",
|
||||
"KnowledgeDecision",
|
||||
"OpenClawCandidateInput",
|
||||
"OpenClawDeliveryPayload",
|
||||
@@ -52,4 +57,4 @@ __all__ = [
|
||||
"build_openclaw_candidate_input",
|
||||
"candidate_id_for",
|
||||
"normalize_candidate_url",
|
||||
]
|
||||
]
|
||||
|
||||
@@ -0,0 +1,35 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import date, datetime
|
||||
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
|
||||
class DailyKeywordTerm(BaseModel):
|
||||
term: str
|
||||
normalized_term: str
|
||||
count: int = Field(ge=1)
|
||||
|
||||
|
||||
class DailyKeywordIndex(BaseModel):
|
||||
schema_version: str = "v1"
|
||||
date: date
|
||||
source: str
|
||||
digest_id: str
|
||||
generated_at: datetime
|
||||
candidate_count: int = Field(default=0, ge=0)
|
||||
terms: list[DailyKeywordTerm] = Field(default_factory=list)
|
||||
|
||||
|
||||
class KeywordStat(BaseModel):
|
||||
term: str
|
||||
total_count: int = Field(default=0, ge=0)
|
||||
days_seen: int = Field(default=0, ge=0)
|
||||
first_seen: date
|
||||
last_seen: date
|
||||
|
||||
|
||||
class KeywordStatsIndex(BaseModel):
|
||||
schema_version: str = "v1"
|
||||
generated_at: datetime
|
||||
terms: list[KeywordStat] = Field(default_factory=list)
|
||||
@@ -6,6 +6,7 @@ from datetime import UTC, date, datetime
|
||||
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
|
||||
@@ -26,8 +27,13 @@ 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:
|
||||
@@ -237,6 +243,17 @@ def run_freshrss_pipeline(
|
||||
)
|
||||
_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,
|
||||
)
|
||||
|
||||
marked_count = 0
|
||||
if mark_read and delivered_item_ids:
|
||||
client.mark_items_as_read(auth_token=auth_token, item_ids=delivered_item_ids)
|
||||
@@ -259,6 +276,7 @@ def run_freshrss_pipeline(
|
||||
"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,
|
||||
}
|
||||
@@ -270,6 +288,7 @@ def run_freshrss_pipeline(
|
||||
"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,
|
||||
|
||||
Reference in New Issue
Block a user