From ae84beeca99ae27c5c54926691ec046fc336c124 Mon Sep 17 00:00:00 2001 From: Bretton Date: Mon, 6 Apr 2026 22:28:08 -0700 Subject: [PATCH] feat(kagi-news): add semantic deduplication via Claude Haiku Adds an LLM-based dedup layer to catch stories that share an event but differ in GUID/headline (same report from different outlets, rewrites, etc.). Exact-GUID filtering still runs first; semantic dedup only considers stories within the same feed posted in a configurable lookback window. Changes: - New SemanticDeduplicator using claude-haiku-4-5 with a tool-use schema for structured duplicate reports; fails open on transient Anthropic API errors so posting is never blocked by network issues - DedupConfig (semantic_enabled, similarity_threshold, lookback_days) with validation, wired through ConfigLoader and config.example.yaml - Aggregator reworked into three phases: GUID filter -> semantic filter -> post, with per-feed counters for new/failed/exact/semantic - StateManager.mark_posted now stores title + summary_snippet, and get_recent_stories returns the lookback window for comparison; state writes are now atomic (tempfile + os.replace) - Requires ANTHROPIC_API_KEY when semantic_enabled is true - Tests: unit tests for SemanticDeduplicator, state manager recent stories + atomic write, and main pipeline integration; live test suite gated behind a `live` pytest marker (excluded by default) Co-Authored-By: Claude Opus 4.6 (1M context) --- aggregators/kagi-news/config.example.yaml | 7 + aggregators/kagi-news/pytest.ini | 3 + aggregators/kagi-news/requirements.txt | 1 + aggregators/kagi-news/src/config.py | 12 +- aggregators/kagi-news/src/main.py | 115 +++++-- aggregators/kagi-news/src/models.py | 19 ++ aggregators/kagi-news/src/semantic_dedup.py | 186 ++++++++++ aggregators/kagi-news/src/state_manager.py | 67 +++- aggregators/kagi-news/tests/test_main.py | 166 ++++++++- .../kagi-news/tests/test_semantic_dedup.py | 249 ++++++++++++++ .../tests/test_semantic_dedup_live.py | 319 ++++++++++++++++++ .../kagi-news/tests/test_state_manager.py | 101 ++++++ 12 files changed, 1204 insertions(+), 41 deletions(-) create mode 100644 aggregators/kagi-news/src/semantic_dedup.py create mode 100644 aggregators/kagi-news/tests/test_semantic_dedup.py create mode 100644 aggregators/kagi-news/tests/test_semantic_dedup_live.py diff --git a/aggregators/kagi-news/config.example.yaml b/aggregators/kagi-news/config.example.yaml index a0415fb..975644b 100644 --- a/aggregators/kagi-news/config.example.yaml +++ b/aggregators/kagi-news/config.example.yaml @@ -28,5 +28,12 @@ feeds: community_handle: "c-science.coves.social" enabled: true +# Semantic deduplication (uses Claude Haiku to detect similar stories) +# Requires ANTHROPIC_API_KEY environment variable when enabled +dedup: + semantic_enabled: true + similarity_threshold: 0.8 # 0.0-1.0, higher = stricter matching + lookback_days: 4 # Compare against stories posted within this window + # Logging configuration log_level: "info" # debug, info, warning, error diff --git a/aggregators/kagi-news/pytest.ini b/aggregators/kagi-news/pytest.ini index 378ac53..301dd2e 100644 --- a/aggregators/kagi-news/pytest.ini +++ b/aggregators/kagi-news/pytest.ini @@ -3,6 +3,8 @@ testpaths = tests python_files = test_*.py python_classes = Test* python_functions = test_* +markers = + live: tests that hit real external APIs (Anthropic) addopts = -v --strict-markers @@ -10,3 +12,4 @@ addopts = --cov=src --cov-report=term-missing --cov-report=html + -m "not live" diff --git a/aggregators/kagi-news/requirements.txt b/aggregators/kagi-news/requirements.txt index 562a50d..88365bd 100644 --- a/aggregators/kagi-news/requirements.txt +++ b/aggregators/kagi-news/requirements.txt @@ -3,6 +3,7 @@ feedparser==6.0.11 beautifulsoup4==4.12.3 requests==2.31.0 pyyaml==6.0.1 +anthropic>=0.40.0 # Testing pytest==8.1.1 diff --git a/aggregators/kagi-news/src/config.py b/aggregators/kagi-news/src/config.py index 0ca2b3f..bb1dabe 100644 --- a/aggregators/kagi-news/src/config.py +++ b/aggregators/kagi-news/src/config.py @@ -10,7 +10,7 @@ from typing import Dict, Any import yaml from urllib.parse import urlparse -from src.models import AggregatorConfig, FeedConfig +from src.models import AggregatorConfig, FeedConfig, DedupConfig logger = logging.getLogger(__name__) @@ -105,12 +105,20 @@ class ConfigLoader: feed = self._parse_feed(feed_data) feeds.append(feed) + # Parse dedup config (optional, defaults if missing) + dedup_data = data.get('dedup', {}) + # Filter to only known DedupConfig fields to avoid TypeError on unknown keys + known_dedup_fields = {f.name for f in DedupConfig.__dataclass_fields__.values()} + filtered_dedup_data = {k: v for k, v in dedup_data.items() if k in known_dedup_fields} + dedup = DedupConfig(**filtered_dedup_data) + logger.info(f"Loaded configuration with {len(feeds)} feeds ({sum(1 for f in feeds if f.enabled)} enabled)") return AggregatorConfig( coves_api_url=coves_api_url, feeds=feeds, - log_level=log_level + log_level=log_level, + dedup=dedup ) def _parse_feed(self, data: Dict[str, Any]) -> FeedConfig: diff --git a/aggregators/kagi-news/src/main.py b/aggregators/kagi-news/src/main.py index ce415d7..b49ca3a 100644 --- a/aggregators/kagi-news/src/main.py +++ b/aggregators/kagi-news/src/main.py @@ -24,6 +24,7 @@ from src.html_parser import KagiHTMLParser from src.richtext_formatter import RichTextFormatter from src.state_manager import StateManager from src.coves_client import CovesClient +from src.semantic_dedup import SemanticDeduplicator # Setup logging logging.basicConfig( @@ -44,7 +45,8 @@ class Aggregator: self, config_path: Path, state_file: Path, - coves_client: Optional[CovesClient] = None + coves_client: Optional[CovesClient] = None, + semantic_dedup: Optional[SemanticDeduplicator] = None ): """ Initialize aggregator. @@ -53,6 +55,7 @@ class Aggregator: config_path: Path to config.yaml state_file: Path to state.json coves_client: Optional CovesClient (for testing) + semantic_dedup: Optional SemanticDeduplicator (for testing) """ # Load configuration logger.info("Loading configuration...") @@ -84,6 +87,25 @@ class Aggregator: api_key=api_key ) + # Initialize semantic deduplicator (or use provided one for testing) + if semantic_dedup: + self.semantic_dedup = semantic_dedup + elif self.config.dedup.semantic_enabled: + anthropic_key = os.getenv('ANTHROPIC_API_KEY') + if not anthropic_key: + raise ValueError( + "ANTHROPIC_API_KEY environment variable required when " + "dedup.semantic_enabled is true. Set the env var or " + "set dedup.semantic_enabled: false in config." + ) + self.semantic_dedup = SemanticDeduplicator( + api_key=anthropic_key, + threshold=self.config.dedup.similarity_threshold + ) + logger.info("Semantic deduplication enabled") + else: + self.semantic_dedup = None + def run(self): """ Run aggregator: fetch, parse, post, and update state. @@ -123,6 +145,11 @@ class Aggregator: """ Process a single RSS feed. + Three phases: + 1. Parse all entries, filter by exact GUID match + 2. Filter by semantic similarity (if enabled) + 3. Post remaining candidates + Args: feed_config: FeedConfig object """ @@ -139,20 +166,19 @@ class Aggregator: if feed.bozo: logger.warning(f"Feed '{feed_config.name}' has parsing issues (bozo flag set)") - # Process entries - new_posts = 0 - skipped_posts = 0 + # Phase 1: Parse all entries, filter by exact GUID + # Store as (entry_guid, story) tuples to preserve the authoritative GUID + candidates = [] + skipped_guid = 0 for entry in feed.entries: try: - # Check if already posted guid = entry.guid if hasattr(entry, 'guid') else entry.link if self.state_manager.is_posted(feed_config.url, guid): - skipped_posts += 1 + skipped_guid += 1 logger.debug(f"Skipping already-posted story: {guid}") continue - # Parse story story = self.html_parser.parse_to_story( title=entry.title, link=entry.link, @@ -161,11 +187,37 @@ class Aggregator: categories=[tag.term for tag in entry.tags] if hasattr(entry, 'tags') else [], html_description=entry.description ) + candidates.append((guid, story)) + + except Exception as e: + logger.error(f"Error processing entry: {e}", exc_info=True) + continue - # Format as rich text + # Phase 2: Semantic dedup (within same feed only) + skipped_semantic = 0 + if self.semantic_dedup and candidates: + recent_stories = self.state_manager.get_recent_stories( + feed_config.url, self.config.dedup.lookback_days + ) + if recent_stories: + new_for_comparison = [ + {"id": guid, "title": story.title, "summary": (story.summary or "")[:200]} + for guid, story in candidates + ] + duplicate_ids = self.semantic_dedup.find_duplicates( + new_for_comparison, recent_stories + ) + before_count = len(candidates) + candidates = [(g, s) for g, s in candidates if g not in duplicate_ids] + skipped_semantic = before_count - len(candidates) + + # Phase 3: Post remaining candidates + new_posts = 0 + failed_posts = 0 + for guid, story in candidates: + try: rich_text = self.richtext_formatter.format_full(story) - # Create external embed with sources sources = [ {"uri": s.url, "title": s.title, "domain": s.domain} for s in story.sources @@ -178,38 +230,35 @@ class Aggregator: sources=sources ) - # Post to community - # Pass thumbnail URL from RSS feed at top level for trusted aggregator upload - try: - post_uri = self.coves_client.create_post( - community_handle=feed_config.community_handle, - title=story.title, - content=rich_text["content"], - facets=rich_text["facets"], - embed=embed, - thumbnail_url=story.image_url # From RSS feed - server will validate and upload - ) - - # Mark as posted (only if successful) - self.state_manager.mark_posted(feed_config.url, guid, post_uri) - new_posts += 1 - logger.info(f"Posted: {story.title[:50]}... -> {post_uri}") - - except Exception as e: - # Don't update state if posting failed - logger.error(f"Failed to post story '{story.title}': {e}") - continue + post_uri = self.coves_client.create_post( + community_handle=feed_config.community_handle, + title=story.title, + content=rich_text["content"], + facets=rich_text["facets"], + embed=embed, + thumbnail_url=story.image_url + ) + + self.state_manager.mark_posted( + feed_config.url, guid, post_uri, + title=story.title, + summary_snippet=story.summary[:200] + ) + new_posts += 1 + logger.info(f"Posted: {story.title[:50]}... -> {post_uri}") except Exception as e: - # Log error but continue with other entries - logger.error(f"Error processing entry: {e}", exc_info=True) + failed_posts += 1 + logger.error(f"Failed to post story '{story.title}' in feed '{feed_config.name}': {e}") continue # Update last run timestamp self.state_manager.update_last_run(feed_config.url, datetime.now()) logger.info( - f"Feed '{feed_config.name}': {new_posts} new posts, {skipped_posts} duplicates" + f"Feed '{feed_config.name}': {new_posts} new, " + f"{failed_posts} failed, " + f"{skipped_guid} exact dupes, {skipped_semantic} semantic dupes" ) diff --git a/aggregators/kagi-news/src/models.py b/aggregators/kagi-news/src/models.py index c8b3d39..8d15129 100644 --- a/aggregators/kagi-news/src/models.py +++ b/aggregators/kagi-news/src/models.py @@ -72,9 +72,28 @@ class FeedConfig: enabled: bool = True +@dataclass(frozen=True) +class DedupConfig: + """Configuration for semantic deduplication.""" + semantic_enabled: bool = True + similarity_threshold: float = 0.8 + lookback_days: int = 4 + + def __post_init__(self): + if not (0.0 <= self.similarity_threshold <= 1.0): + raise ValueError( + f"similarity_threshold must be between 0.0 and 1.0, got {self.similarity_threshold}" + ) + if self.lookback_days < 1: + raise ValueError( + f"lookback_days must be >= 1, got {self.lookback_days}" + ) + + @dataclass class AggregatorConfig: """Full aggregator configuration.""" coves_api_url: str feeds: List[FeedConfig] log_level: str = "info" + dedup: DedupConfig = field(default_factory=DedupConfig) diff --git a/aggregators/kagi-news/src/semantic_dedup.py b/aggregators/kagi-news/src/semantic_dedup.py new file mode 100644 index 0000000..821df92 --- /dev/null +++ b/aggregators/kagi-news/src/semantic_dedup.py @@ -0,0 +1,186 @@ +""" +Semantic Deduplication for Kagi News Aggregator. + +Uses Claude Haiku to detect semantically similar stories within the same +feed/community. Compares new candidate stories against recently posted ones +using a single batched API call per feed. +""" +import logging +from typing import List, Dict, Set + +import anthropic + +logger = logging.getLogger(__name__) + +# Tool schema for structured output from Haiku +REPORT_DUPLICATES_TOOL = { + "name": "report_duplicates", + "description": "Report which new articles are semantic duplicates of recent ones", + "input_schema": { + "type": "object", + "properties": { + "results": { + "type": "array", + "items": { + "type": "object", + "properties": { + "new_id": { + "type": "string", + "description": "The ID of the new candidate article" + }, + "duplicate_of": { + "type": "string", + "description": "The ID of the recent article this is a duplicate of, or empty string if not a duplicate" + }, + "confidence": { + "type": "number", + "description": "Confidence score 0.0-1.0 that this is a duplicate" + } + }, + "required": ["new_id", "duplicate_of", "confidence"] + } + } + }, + "required": ["results"] + } +} + +SYSTEM_PROMPT = """You are a news deduplication system. Your job is to identify when a NEW article is essentially a rewrite of the same report as a RECENT article — covering the identical event with no meaningful new information. + +DUPLICATE means: both articles describe the same specific occurrence at the same point in time, and the new one adds no significant new facts, outcome, or development beyond what the other already reported. They are essentially the same report from different outlets or with different headlines. + +CRITICAL: For ongoing stories (missions, conflicts, investigations, court cases), focus on the article's PRIMARY new fact — NOT on recap/background context. News articles routinely recap prior events; that recap does NOT make the new article a duplicate. Compare what the article is actually REPORTING NEW. + +NOT a duplicate when ANY of the following apply: +- The new article reports a SUBSEQUENT event, even if closely related — this includes: + • Anticipation/preparation followed by the actual event ("preparing to launch" vs "launched successfully") + • Announcement followed by implementation ("tariffs announced" vs "tariffs take effect") + • Deadline followed by response ("48-hour ultimatum" vs "Iran responds") + • One milestone followed by the next milestone (orbit → translunar injection → flyby → return) +- The new article reports a new daily status update on an ongoing mission/event, even if it recaps prior days +- The new article adds substantial new information: a new outcome, official response, policy change, or follow-up action +- The articles cover the same broad topic but different specific events or different time points + +Examples: +- DUPLICATE: "US revokes residency of Soleimani relatives, detains two" and "US revokes visas of Soleimani relatives, detains two" (same event rewritten same day) +- DUPLICATE: "Planet Labs halts Iran conflict imagery after US request" and "US satellite firm halts Iran conflict imagery release" (same announcement) +- DUPLICATE: "Italian court orders Netflix to refund subscribers for price hikes" and "Italian court orders Netflix to refund customers up to €500" (same court ruling) +- NOT DUPLICATE: "NASA readies Artemis II crewed moon flyby" and "NASA launches Artemis II crew on lunar flyby" (preparation vs the launch actually happening — DIFFERENT events on different days) +- NOT DUPLICATE: "Artemis II sends astronauts toward the Moon" (TLI burn) and "Artemis II tests moon mission technologies" (next-day status, different activity) +- NOT DUPLICATE: "Artemis II crew nears moon flyby" (halfway point) and "Artemis II crew enters lunar space" (sphere of influence reached) — distinct mission milestones +- NOT DUPLICATE: "Trump gives Iran 48-hour Hormuz deadline" and "Iran allows Iraqi ships through Strait of Hormuz" (sequential events in same crisis) +- NOT DUPLICATE: "US announces tariffs on China" and "Tariffs officially take effect, markets react" (announcement vs implementation) +- NOT DUPLICATE: "Earthquake hits Turkey" and "Earthquake hits Japan" (same topic, different events) + +When in doubt, mark as NOT a duplicate. It is better to allow a near-duplicate through than to suppress a unique story or a daily status update on an ongoing event.""" + + +class SemanticDeduplicator: + """ + Detects semantically similar news stories using Claude Haiku. + + Uses a single batched API call per feed to compare all new candidates + against all recent stories simultaneously. + """ + + def __init__(self, api_key: str, threshold: float = 0.8, + model: str = "claude-haiku-4-5-20251001"): + """ + Initialize the semantic deduplicator. + + Args: + api_key: Anthropic API key + threshold: Minimum confidence to consider a duplicate (0.0-1.0) + model: Claude model to use + """ + self.client = anthropic.Anthropic(api_key=api_key) + self.threshold = threshold + self.model = model + + def find_duplicates( + self, + new_stories: List[Dict], + recent_stories: List[Dict], + ) -> Set[str]: + """ + Find which new stories are semantic duplicates of recent ones. + + Args: + new_stories: List of dicts with keys: id, title, summary + recent_stories: List of dicts with keys: id, title, summary + + Returns: + Set of new story IDs that are semantic duplicates + """ + if not new_stories or not recent_stories: + return set() + + prompt = self._build_prompt(new_stories, recent_stories) + + try: + response = self.client.messages.create( + model=self.model, + max_tokens=1024, + system=SYSTEM_PROMPT, + tools=[REPORT_DUPLICATES_TOOL], + tool_choice={"type": "tool", "name": "report_duplicates"}, + messages=[{"role": "user", "content": prompt}] + ) + + return self._parse_response(response) + + except (anthropic.APIConnectionError, anthropic.RateLimitError, + anthropic.InternalServerError) as e: + # Fail open for transient/network errors: don't block any posts + logger.warning(f"Semantic dedup API call failed (transient), allowing all posts: {e}") + return set() + + def _build_prompt(self, new_stories: List[Dict], + recent_stories: List[Dict]) -> str: + """Build the comparison prompt for Haiku.""" + lines = ["Compare each NEW article against the RECENT articles and identify duplicates.\n"] + + lines.append("RECENT articles (already posted):") + for story in recent_stories: + summary_part = f" -- {story['summary']}" if story.get('summary') else "" + lines.append(f"- [{story['id']}] {story['title']}{summary_part}") + + lines.append("\nNEW candidates:") + for story in new_stories: + summary_part = f" -- {story['summary']}" if story.get('summary') else "" + lines.append(f"- [{story['id']}] {story['title']}{summary_part}") + + lines.append(f"\nUse the report_duplicates tool. Only mark as duplicate if confidence >= {self.threshold}.") + + return "\n".join(lines) + + def _parse_response(self, response) -> Set[str]: + """Parse the tool use response and return duplicate story IDs.""" + duplicates = set() + + for block in response.content: + if block.type == "tool_use" and block.name == "report_duplicates": + results = block.input.get("results", []) + for result in results: + duplicate_of = result.get("duplicate_of", "") + confidence = result.get("confidence", 0.0) + + new_id = result.get("new_id") + if not new_id: + logger.warning(f"Skipping dedup result missing 'new_id': {result}") + continue + + if duplicate_of and confidence >= self.threshold: + duplicates.add(new_id) + logger.info( + f"Semantic duplicate detected: '{new_id}' " + f"duplicates '{duplicate_of}' " + f"(confidence: {confidence:.2f})" + ) + + if duplicates: + logger.info(f"Semantic dedup filtered {len(duplicates)} duplicate(s)") + else: + logger.info("Semantic dedup: no duplicates found") + + return duplicates diff --git a/aggregators/kagi-news/src/state_manager.py b/aggregators/kagi-news/src/state_manager.py index 9063ef4..5ecaaa8 100644 --- a/aggregators/kagi-news/src/state_manager.py +++ b/aggregators/kagi-news/src/state_manager.py @@ -6,6 +6,8 @@ Uses JSON file for persistence. """ import json import logging +import os +import tempfile from pathlib import Path from datetime import datetime, timedelta from typing import Optional, Dict, List @@ -64,8 +66,21 @@ class StateManager: # Ensure parent directory exists self.state_file.parent.mkdir(parents=True, exist_ok=True) - with open(self.state_file, 'w') as f: - json.dump(state, f, indent=2) + # Atomic write: write to temp file then rename to avoid corruption + fd, tmp_path = tempfile.mkstemp( + dir=str(self.state_file.parent), suffix='.tmp' + ) + try: + with os.fdopen(fd, 'w') as f: + json.dump(state, f, indent=2) + os.replace(tmp_path, str(self.state_file)) + except BaseException: + # Clean up temp file on any failure + try: + os.unlink(tmp_path) + except OSError: + pass + raise def _ensure_feed_exists(self, feed_url: str): """Ensure feed entry exists in state.""" @@ -91,7 +106,8 @@ class StateManager: posted_guids = self.state['feeds'][feed_url]['posted_guids'] return any(entry['guid'] == guid for entry in posted_guids) - def mark_posted(self, feed_url: str, guid: str, post_uri: str): + def mark_posted(self, feed_url: str, guid: str, post_uri: str, + title: str = "", summary_snippet: str = ""): """ Mark a story as posted. @@ -99,6 +115,8 @@ class StateManager: feed_url: RSS feed URL guid: Story GUID post_uri: AT Proto URI of created post + title: Story title (for semantic dedup) + summary_snippet: First 200 chars of summary (for semantic dedup) """ self._ensure_feed_exists(feed_url) @@ -106,7 +124,9 @@ class StateManager: entry = { 'guid': guid, 'post_uri': post_uri, - 'posted_at': datetime.now().isoformat() + 'posted_at': datetime.now().isoformat(), + 'title': title, + 'summary_snippet': summary_snippet[:200] if summary_snippet else "" } self.state['feeds'][feed_url]['posted_guids'].append(entry) @@ -186,6 +206,45 @@ class StateManager: if old_count != new_count: logger.info(f"Cleaned up {old_count - new_count} old entries for {feed_url}") + def get_recent_stories(self, feed_url: str, days: int = 4) -> List[Dict]: + """ + Get recently posted stories with title and summary for semantic comparison. + + Args: + feed_url: RSS feed URL + days: Number of days to look back (default: 4) + + Returns: + List of dicts with keys: id, title, summary + Only includes entries that have title data (backward compatible). + """ + self._ensure_feed_exists(feed_url) + + cutoff = datetime.now() - timedelta(days=days) + recent = [] + + for entry in self.state['feeds'][feed_url]['posted_guids']: + try: + # Skip entries without title (old format) + title = entry.get('title', '') + if not title: + continue + + posted_at = datetime.fromisoformat(entry['posted_at']) + if posted_at > cutoff: + recent.append({ + 'id': entry['guid'], + 'title': title, + 'summary': entry.get('summary_snippet', '') + }) + except (ValueError, KeyError) as e: + logger.warning( + f"Skipping malformed state entry for feed '{feed_url}': {e}" + ) + continue + + return recent + def get_posted_count(self, feed_url: str) -> int: """ Get count of posted items for a feed. diff --git a/aggregators/kagi-news/tests/test_main.py b/aggregators/kagi-news/tests/test_main.py index af67c40..d34cf50 100644 --- a/aggregators/kagi-news/tests/test_main.py +++ b/aggregators/kagi-news/tests/test_main.py @@ -10,7 +10,7 @@ from unittest.mock import Mock, MagicMock, patch, call import feedparser from src.main import Aggregator -from src.models import KagiStory, AggregatorConfig, FeedConfig, Perspective, Quote, Source +from src.models import KagiStory, AggregatorConfig, FeedConfig, DedupConfig, Perspective, Quote, Source @pytest.fixture @@ -38,7 +38,8 @@ def mock_config(): enabled=False ) ], - log_level="info" + log_level="info", + dedup=DedupConfig(semantic_enabled=False) ) @@ -596,3 +597,164 @@ class TestAggregator: # Verify sources is None (empty list becomes None) assert call_kwargs.get("sources") is None + + def test_semantic_dedup_filters_duplicates(self, mock_rss_feed, tmp_path): + """Test that semantic dedup filters out similar stories when recent stories exist.""" + import json + + # Config with semantic dedup enabled + config = AggregatorConfig( + coves_api_url="https://api.coves.social", + feeds=[ + FeedConfig( + name="World News", + url="https://news.kagi.com/world.xml", + community_handle="world-news.coves.social", + enabled=True + ) + ], + log_level="info", + dedup=DedupConfig(semantic_enabled=True, lookback_days=4) + ) + + # Pre-populate state with recent stories so get_recent_stories returns data + state_file = tmp_path / "state.json" + state_file.write_text(json.dumps({ + "feeds": { + "https://news.kagi.com/world.xml": { + "posted_guids": [ + { + "guid": "existing-1", + "post_uri": "at://test/1", + "posted_at": datetime.now().isoformat(), + "title": "US announces new tariffs on China", + "summary_snippet": "The United States has announced new tariffs." + } + ], + "last_successful_run": None + } + } + })) + + mock_client = Mock() + mock_client.create_post.return_value = "at://did:plc:test/social.coves.post/abc123" + + # Create two distinct stories so we can verify one is filtered and one is not + story_1 = KagiStory( + title="Trade tensions escalate as US tariffs take effect", + link="https://kite.kagi.com/test/world/1", + guid="https://kite.kagi.com/test/world/1", + pub_date=datetime(2024, 1, 15, 12, 0, 0), + categories=["World"], + summary="US tariffs on Chinese goods go into effect.", + highlights=[], perspectives=[], quote=None, sources=[], + image_url=None, image_alt=None + ) + story_2 = KagiStory( + title="Earthquake hits Turkey killing dozens", + link="https://kite.kagi.com/test/world/2", + guid="https://kite.kagi.com/test/world/2", + pub_date=datetime(2024, 1, 15, 13, 0, 0), + categories=["World"], + summary="A 6.5 magnitude earthquake struck southeastern Turkey.", + highlights=[], perspectives=[], quote=None, sources=[], + image_url=None, image_alt=None + ) + + # Mock semantic dedup to mark first story as duplicate of existing-1 + mock_dedup = Mock() + mock_dedup.find_duplicates.return_value = {"https://kite.kagi.com/test/world/1"} + + with patch('src.main.ConfigLoader') as MockConfigLoader, \ + patch('src.main.RSSFetcher') as MockRSSFetcher, \ + patch('src.main.KagiHTMLParser') as MockHTMLParser, \ + patch('src.main.RichTextFormatter') as MockFormatter: + + mock_loader = Mock() + mock_loader.load.return_value = config + MockConfigLoader.return_value = mock_loader + + mock_fetcher = Mock() + mock_fetcher.fetch_feed.return_value = mock_rss_feed + MockRSSFetcher.return_value = mock_fetcher + + mock_parser = Mock() + # Return different stories for the two entries + mock_parser.parse_to_story.side_effect = [story_1, story_2] + MockHTMLParser.return_value = mock_parser + + mock_formatter = Mock() + mock_formatter.format_full.return_value = { + "content": "Test content", + "facets": [] + } + MockFormatter.return_value = mock_formatter + + aggregator = Aggregator( + config_path=Path("config.yaml"), + state_file=state_file, + coves_client=mock_client, + semantic_dedup=mock_dedup + ) + aggregator.run() + + # find_duplicates should have been called with the new stories + mock_dedup.find_duplicates.assert_called_once() + call_args = mock_dedup.find_duplicates.call_args + new_for_comparison = call_args[0][0] + recent_for_comparison = call_args[0][1] + + # Should have passed both new candidates + assert len(new_for_comparison) == 2 + # Should have passed the pre-populated recent story + assert len(recent_for_comparison) == 1 + assert recent_for_comparison[0]["id"] == "existing-1" + + # Only story_2 should be posted (story_1 was marked as duplicate) + assert mock_client.create_post.call_count == 1 + posted_title = mock_client.create_post.call_args.kwargs.get("title") + assert posted_title == "Earthquake hits Turkey killing dozens" + + def test_semantic_dedup_disabled_skips_check(self, mock_config, mock_rss_feed, sample_story, tmp_path): + """Test that semantic dedup is skipped when disabled.""" + state_file = tmp_path / "state.json" + mock_client = Mock() + mock_client.create_post.return_value = "at://did:plc:test/social.coves.post/abc123" + + mock_dedup = Mock() + + with patch('src.main.ConfigLoader') as MockConfigLoader, \ + patch('src.main.RSSFetcher') as MockRSSFetcher, \ + patch('src.main.KagiHTMLParser') as MockHTMLParser, \ + patch('src.main.RichTextFormatter') as MockFormatter: + + mock_loader = Mock() + mock_loader.load.return_value = mock_config # dedup disabled + MockConfigLoader.return_value = mock_loader + + mock_fetcher = Mock() + mock_fetcher.fetch_feed.return_value = mock_rss_feed + MockRSSFetcher.return_value = mock_fetcher + + mock_parser = Mock() + mock_parser.parse_to_story.return_value = sample_story + MockHTMLParser.return_value = mock_parser + + mock_formatter = Mock() + mock_formatter.format_full.return_value = { + "content": "Test content", + "facets": [] + } + MockFormatter.return_value = mock_formatter + + aggregator = Aggregator( + config_path=Path("config.yaml"), + state_file=state_file, + coves_client=mock_client + ) + # semantic_dedup should be None when disabled + assert aggregator.semantic_dedup is None + + aggregator.run() + # All stories should be posted (no semantic filtering) + assert mock_client.create_post.call_count == 4 diff --git a/aggregators/kagi-news/tests/test_semantic_dedup.py b/aggregators/kagi-news/tests/test_semantic_dedup.py new file mode 100644 index 0000000..29471de --- /dev/null +++ b/aggregators/kagi-news/tests/test_semantic_dedup.py @@ -0,0 +1,249 @@ +""" +Tests for Semantic Deduplication module. + +Tests the SemanticDeduplicator with mocked Anthropic API responses. +""" +import pytest +from unittest.mock import Mock, MagicMock, patch +import httpx + +import anthropic +from src.semantic_dedup import SemanticDeduplicator, REPORT_DUPLICATES_TOOL + + +@pytest.fixture +def dedup(): + """Create a SemanticDeduplicator with mocked client.""" + with patch('src.semantic_dedup.anthropic.Anthropic') as mock_anthropic_cls: + mock_client = Mock() + mock_anthropic_cls.return_value = mock_client + d = SemanticDeduplicator(api_key="test-key", threshold=0.8) + d.client = mock_client + yield d + + +@pytest.fixture +def recent_stories(): + """Sample recent stories for comparison.""" + return [ + {"id": "recent-1", "title": "US announces new tariffs on China", "summary": "The United States has announced sweeping new tariffs on Chinese goods."}, + {"id": "recent-2", "title": "SpaceX launches Starship for 5th test flight", "summary": "SpaceX successfully launched Starship on its fifth test flight."}, + {"id": "recent-3", "title": "New study finds coffee reduces heart disease risk", "summary": "A comprehensive study shows moderate coffee consumption lowers cardiovascular risk."}, + ] + + +@pytest.fixture +def new_stories(): + """Sample new candidate stories.""" + return [ + {"id": "new-1", "title": "Trade tensions escalate as US tariffs take effect", "summary": "US tariffs on Chinese goods go into effect amid growing trade war."}, + {"id": "new-2", "title": "Earthquake hits Turkey killing dozens", "summary": "A 6.5 magnitude earthquake struck southeastern Turkey."}, + {"id": "new-3", "title": "EU considers new trade deal with Japan", "summary": "European Union opens talks on a new bilateral trade agreement with Japan."}, + ] + + +def _make_tool_response(results): + """Helper to create a mock Anthropic tool use response.""" + mock_response = Mock() + mock_block = Mock() + mock_block.type = "tool_use" + mock_block.name = "report_duplicates" + mock_block.input = {"results": results} + mock_response.content = [mock_block] + return mock_response + + +class TestSemanticDeduplicator: + """Test suite for SemanticDeduplicator.""" + + def test_find_duplicates_detects_similar_story(self, dedup, new_stories, recent_stories): + """Test that semantically similar stories are detected.""" + dedup.client.messages.create.return_value = _make_tool_response([ + {"new_id": "new-1", "duplicate_of": "recent-1", "confidence": 0.92}, + {"new_id": "new-2", "duplicate_of": "", "confidence": 0.0}, + {"new_id": "new-3", "duplicate_of": "", "confidence": 0.0}, + ]) + + duplicates = dedup.find_duplicates(new_stories, recent_stories) + + assert duplicates == {"new-1"} + dedup.client.messages.create.assert_called_once() + + def test_find_duplicates_no_matches(self, dedup, new_stories, recent_stories): + """Test when no duplicates are found.""" + dedup.client.messages.create.return_value = _make_tool_response([ + {"new_id": "new-1", "duplicate_of": "", "confidence": 0.0}, + {"new_id": "new-2", "duplicate_of": "", "confidence": 0.0}, + {"new_id": "new-3", "duplicate_of": "", "confidence": 0.0}, + ]) + + duplicates = dedup.find_duplicates(new_stories, recent_stories) + + assert duplicates == set() + + def test_threshold_filtering(self, dedup, new_stories, recent_stories): + """Test that confidence below threshold is not filtered.""" + dedup.client.messages.create.return_value = _make_tool_response([ + {"new_id": "new-1", "duplicate_of": "recent-1", "confidence": 0.7}, # Below 0.8 threshold + {"new_id": "new-2", "duplicate_of": "", "confidence": 0.0}, + {"new_id": "new-3", "duplicate_of": "recent-1", "confidence": 0.85}, # Above threshold + ]) + + duplicates = dedup.find_duplicates(new_stories, recent_stories) + + assert "new-1" not in duplicates # Below threshold + assert "new-3" in duplicates # Above threshold + + def test_empty_new_stories_no_api_call(self, dedup, recent_stories): + """Test that empty new stories list skips API call.""" + duplicates = dedup.find_duplicates([], recent_stories) + + assert duplicates == set() + dedup.client.messages.create.assert_not_called() + + def test_empty_recent_stories_no_api_call(self, dedup, new_stories): + """Test that empty recent stories list skips API call.""" + duplicates = dedup.find_duplicates(new_stories, []) + + assert duplicates == set() + dedup.client.messages.create.assert_not_called() + + def test_fail_open_on_transient_error(self, dedup, new_stories, recent_stories): + """Test that transient API errors result in no filtering (fail open).""" + dedup.client.messages.create.side_effect = anthropic.APIConnectionError( + request=httpx.Request("POST", "https://api.anthropic.com") + ) + + duplicates = dedup.find_duplicates(new_stories, recent_stories) + + assert duplicates == set() + + def test_auth_error_propagates(self, dedup, new_stories, recent_stories): + """Test that authentication errors propagate instead of being silently swallowed.""" + mock_response = httpx.Response(401, request=httpx.Request("POST", "https://api.anthropic.com")) + dedup.client.messages.create.side_effect = anthropic.AuthenticationError( + message="Invalid API key", + response=mock_response, + body={"error": {"message": "Invalid API key"}}, + ) + + with pytest.raises(anthropic.AuthenticationError): + dedup.find_duplicates(new_stories, recent_stories) + + def test_multiple_duplicates_detected(self, dedup, recent_stories): + """Test detecting multiple duplicates in one batch.""" + new_stories = [ + {"id": "new-1", "title": "Tariffs take effect on Chinese imports", "summary": "US tariffs begin."}, + {"id": "new-2", "title": "Starship fifth flight declared success", "summary": "SpaceX celebrates."}, + {"id": "new-3", "title": "Unique story about Mars", "summary": "Mars exploration update."}, + ] + + dedup.client.messages.create.return_value = _make_tool_response([ + {"new_id": "new-1", "duplicate_of": "recent-1", "confidence": 0.90}, + {"new_id": "new-2", "duplicate_of": "recent-2", "confidence": 0.88}, + {"new_id": "new-3", "duplicate_of": "", "confidence": 0.0}, + ]) + + duplicates = dedup.find_duplicates(new_stories, recent_stories) + + assert duplicates == {"new-1", "new-2"} + + def test_uses_tool_use_for_structured_output(self, dedup, new_stories, recent_stories): + """Test that the API call uses tool_choice to force structured output.""" + dedup.client.messages.create.return_value = _make_tool_response([ + {"new_id": "new-1", "duplicate_of": "", "confidence": 0.0}, + {"new_id": "new-2", "duplicate_of": "", "confidence": 0.0}, + {"new_id": "new-3", "duplicate_of": "", "confidence": 0.0}, + ]) + + dedup.find_duplicates(new_stories, recent_stories) + + call_kwargs = dedup.client.messages.create.call_args.kwargs + assert call_kwargs["tool_choice"] == {"type": "tool", "name": "report_duplicates"} + assert call_kwargs["tools"] == [REPORT_DUPLICATES_TOOL] + + def test_build_prompt_includes_all_stories(self, dedup): + """Test that the prompt includes all recent and new stories.""" + new = [{"id": "n1", "title": "New Title", "summary": "New summary"}] + recent = [{"id": "r1", "title": "Recent Title", "summary": "Recent summary"}] + + prompt = dedup._build_prompt(new, recent) + + assert "New Title" in prompt + assert "New summary" in prompt + assert "Recent Title" in prompt + assert "Recent summary" in prompt + assert "[n1]" in prompt + assert "[r1]" in prompt + + def test_build_prompt_handles_empty_summary(self, dedup): + """Test prompt building with empty summaries.""" + new = [{"id": "n1", "title": "New Title", "summary": ""}] + recent = [{"id": "r1", "title": "Recent Title", "summary": ""}] + + prompt = dedup._build_prompt(new, recent) + + assert "New Title" in prompt + assert "Recent Title" in prompt + # Should not have dangling " -- " for empty summaries + assert "-- \n" not in prompt + + def test_parse_response_no_tool_use_blocks(self, dedup): + """Test _parse_response returns empty set when response has no tool_use blocks.""" + mock_response = Mock() + mock_text_block = Mock() + mock_text_block.type = "text" + mock_text_block.text = "Here are some duplicates I found." + mock_response.content = [mock_text_block] + + result = dedup._parse_response(mock_response) + + assert result == set() + + def test_parse_response_wrong_tool_name(self, dedup): + """Test _parse_response returns empty set when tool_use block has wrong name.""" + mock_response = Mock() + mock_block = Mock() + mock_block.type = "tool_use" + mock_block.name = "wrong_tool_name" + mock_block.input = {"results": [ + {"new_id": "new-1", "duplicate_of": "recent-1", "confidence": 0.95} + ]} + mock_response.content = [mock_block] + + result = dedup._parse_response(mock_response) + + assert result == set() + + def test_threshold_boundary_exactly_equal(self, dedup, new_stories, recent_stories): + """Test that confidence exactly equal to threshold (0.8) IS filtered (>= behavior).""" + dedup.client.messages.create.return_value = _make_tool_response([ + {"new_id": "new-1", "duplicate_of": "recent-1", "confidence": 0.8}, # Exactly at 0.8 threshold + {"new_id": "new-2", "duplicate_of": "", "confidence": 0.0}, + {"new_id": "new-3", "duplicate_of": "", "confidence": 0.0}, + ]) + + duplicates = dedup.find_duplicates(new_stories, recent_stories) + + # 0.8 >= 0.8 threshold, so new-1 SHOULD be filtered + assert "new-1" in duplicates + + def test_custom_threshold(self): + """Test that custom threshold is respected.""" + with patch('src.semantic_dedup.anthropic.Anthropic') as mock_anthropic_cls: + mock_client = Mock() + mock_anthropic_cls.return_value = mock_client + d = SemanticDeduplicator(api_key="test-key", threshold=0.95) + d.client = mock_client + + d.client.messages.create.return_value = _make_tool_response([ + {"new_id": "n1", "duplicate_of": "r1", "confidence": 0.90}, + ]) + + new = [{"id": "n1", "title": "Title", "summary": "Summary"}] + recent = [{"id": "r1", "title": "Title", "summary": "Summary"}] + + duplicates = d.find_duplicates(new, recent) + + # 0.90 < 0.95 threshold, should NOT be filtered + assert duplicates == set() diff --git a/aggregators/kagi-news/tests/test_semantic_dedup_live.py b/aggregators/kagi-news/tests/test_semantic_dedup_live.py new file mode 100644 index 0000000..5d0c07d --- /dev/null +++ b/aggregators/kagi-news/tests/test_semantic_dedup_live.py @@ -0,0 +1,319 @@ +""" +Live integration tests for Semantic Deduplication. + +These tests hit the real Anthropic API with Claude Haiku to validate +that semantic dedup correctly distinguishes: +1. True duplicates (same event, different headline) → SHOULD be filtered +2. Updates/developments (new info on same topic) → should NOT be filtered + +Run with: ANTHROPIC_API_KEY=sk-ant-... pytest tests/test_semantic_dedup_live.py -v -s +Skip in CI with: pytest -m "not live" +""" +import os +import pytest + +from src.semantic_dedup import SemanticDeduplicator + + +# Skip entire module if no API key +pytestmark = pytest.mark.live +ANTHROPIC_API_KEY = os.getenv("ANTHROPIC_API_KEY") + + +@pytest.fixture +def dedup(): + """Create a SemanticDeduplicator with real API key.""" + if not ANTHROPIC_API_KEY: + pytest.skip("ANTHROPIC_API_KEY not set — skipping live test") + return SemanticDeduplicator(api_key=ANTHROPIC_API_KEY, threshold=0.8) + + +class TestDuplicateDetection: + """True duplicates: same event rewritten — these SHOULD be filtered.""" + + def test_soleimani_visa_revocation_cross_feed(self, dedup): + """Same event (Soleimani relatives detained) in World vs USA feeds.""" + recent = [{ + "id": "recent-world-soleimani", + "title": "US revokes residency of Soleimani relatives, detains two", + "summary": ( + "The US revoked green cards and visas for Iranian nationals " + "tied to Iran's leadership, including Hamideh Soleimani Afshar " + "and her daughter. Secretary of State Marco Rubio determined " + "they were ineligible for lawful permanent resident status." + ), + }] + + new = [{ + "id": "new-usa-soleimani", + "title": "US revokes visas of Soleimani relatives, detains two", + "summary": ( + "The Trump administration revoked the green cards or visas of " + "at least four Iranian nationals tied to Iran's current or " + "former government, including Hamideh Soleimani Afshar and " + "her daughter. The two women were detained by ICE in Southern " + "California." + ), + }] + + duplicates = dedup.find_duplicates(new, recent) + assert "new-usa-soleimani" in duplicates, ( + "Same event (Soleimani visa revocation) should be detected as duplicate" + ) + + def test_planet_labs_imagery_halt_cross_feed(self, dedup): + """Same announcement (Planet Labs halts imagery) in World vs Tech feeds.""" + recent = [{ + "id": "recent-world-planet", + "title": "US satellite firm halts Iran conflict imagery release", + "summary": ( + "Planet Labs announced it will indefinitely withhold public " + "release of satellite images covering Iran and nearby conflict " + "zones following a US government request." + ), + }] + + new = [{ + "id": "new-tech-planet", + "title": "Planet Labs halts Iran conflict imagery after US request", + "summary": ( + "Planet Labs, a US commercial satellite imaging company, will " + "indefinitely withhold imagery of Iran and the wider conflict " + "region following a US government request." + ), + }] + + duplicates = dedup.find_duplicates(new, recent) + assert "new-tech-planet" in duplicates, ( + "Same announcement (Planet Labs) should be detected as duplicate" + ) + + def test_italian_netflix_ruling_real_posts(self, dedup): + """Real cross-day duplicate from c-tech.coves.social (04-05 → 04-06).""" + recent = [{ + "id": "3miqvwsfw6s2x", + "title": "Italian court orders Netflix to refund subscribers for price hikes since 2017", + "summary": ( + "A court in Rome ruled that Netflix must refund Italian " + "customers for price increases imposed between 2017 and " + "January 2024, and cut subscription prices back to earlier levels." + ), + }] + + new = [{ + "id": "3mitgfmxetk2x", + "title": "Italian court orders Netflix to refund customers up to €500 for price hikes", + "summary": ( + "An Italian court has ruled that Netflix must refund customers " + "for price increases introduced between 2017 and 2024, with " + "some subscribers potentially receiving up to €500 in reimbursements." + ), + }] + + duplicates = dedup.find_duplicates(new, recent) + assert "3mitgfmxetk2x" in duplicates, ( + "Real cross-day rewrite of Italian Netflix court ruling should be flagged" + ) + + def test_airman_rescue_cross_feed(self, dedup): + """Same rescue event in World vs USA feeds.""" + recent = [{ + "id": "recent-world-rescue", + "title": "Update: US rescues second airman from downed jet", + "summary": ( + "US forces rescued the second crew member from an F-15E " + "shot down over Iran, ending a tense search in the sixth " + "week of conflict." + ), + }] + + new = [{ + "id": "new-usa-rescue", + "title": "Update: U.S. forces rescue downed airman in Iran", + "summary": ( + "In the latest turn in the U.S.-Iran conflict, President " + "Trump said early Sunday that U.S. forces rescued an Air " + "Force officer whose F-15E Strike Eagle was shot down over " + "Iran on Friday." + ), + }] + + duplicates = dedup.find_duplicates(new, recent) + assert "new-usa-rescue" in duplicates, ( + "Same rescue event should be detected as duplicate" + ) + + +class TestUpdatePassthrough: + """Updates/developments: new info on same topic — should NOT be filtered.""" + + def test_hormuz_deadline_vs_response(self, dedup): + """Sequential events: US ultimatum then Iran's response — NOT duplicates.""" + recent = [{ + "id": "recent-hormuz-deadline", + "title": "President Trump gives Iran 48-hour Hormuz deadline", + "summary": ( + "President Trump said on April 4 that Iran had 48 hours to " + "make a deal or reopen the Strait of Hormuz, or face further " + "U.S. action." + ), + }] + + new = [{ + "id": "new-hormuz-response", + "title": "Iran allows Iraqi ships through Strait of Hormuz", + "summary": ( + "Iran announced on April 5 that Iraqi vessels can transit " + "the Strait of Hormuz despite broader restrictions, with " + "military spokesperson describing Iraq as exempt." + ), + }] + + duplicates = dedup.find_duplicates(new, recent) + assert "new-hormuz-response" not in duplicates, ( + "Sequential Hormuz developments should NOT be flagged as duplicate" + ) + + def test_tariff_announcement_vs_implementation(self, dedup): + """Announcement then implementation — NOT duplicates.""" + recent = [{ + "id": "recent-tariff-announce", + "title": "US announces sweeping new tariffs on China", + "summary": ( + "The United States announced new tariffs on Chinese goods " + "effective next week, targeting electronics and automotive parts." + ), + }] + + new = [{ + "id": "new-tariff-effect", + "title": "US-China tariffs take effect, markets drop sharply", + "summary": ( + "The new US tariffs on Chinese goods officially took effect " + "today, sending markets tumbling. The S&P 500 fell 2.3% as " + "Beijing threatened retaliatory measures." + ), + }] + + duplicates = dedup.find_duplicates(new, recent) + assert "new-tariff-effect" not in duplicates, ( + "Tariff implementation (with market reaction) is a new development, not a duplicate" + ) + + def test_samsung_different_products(self, dedup): + """Same company, different product news — NOT duplicates.""" + recent = [{ + "id": "recent-samsung-messages", + "title": "Samsung discontinues Messages app in favor of Google Messages", + "summary": ( + "Samsung will discontinue its Messages app in the United " + "States in July 2026 and transition users to Google Messages." + ), + }] + + new = [{ + "id": "new-samsung-s26", + "title": "Samsung Galaxy S26 Ultra sparks camera upgrade debate", + "summary": ( + "Early coverage of Samsung's Galaxy S26 Ultra presents mixed " + "perspectives on whether a new flagship justifies upgrading." + ), + }] + + duplicates = dedup.find_duplicates(new, recent) + assert "new-samsung-s26" not in duplicates, ( + "Different Samsung product news should NOT be flagged as duplicate" + ) + + def test_climate_different_findings(self, dedup): + """Same broad topic (climate), different research — NOT duplicates.""" + recent = [{ + "id": "recent-climate-spring", + "title": "Climate researchers link earlier spring to warming trends", + "summary": ( + "Spring conditions are arriving approximately seven days " + "earlier in St. Louis according to climate experts." + ), + }] + + new = [{ + "id": "new-climate-arctic", + "title": "Arctic winter sea ice ties record low", + "summary": ( + "Arctic winter sea ice reached its seasonal maximum at " + "record-low levels. Thawing permafrost across northern " + "Alaska increases runoff." + ), + }] + + duplicates = dedup.find_duplicates(new, recent) + assert "new-climate-arctic" not in duplicates, ( + "Different climate research should NOT be flagged as duplicate" + ) + + def test_artemis_different_milestones(self, dedup): + """Same mission, different milestones — NOT duplicates.""" + recent = [{ + "id": "recent-artemis-launch", + "title": "NASA Artemis II crew reaches Earth orbit successfully", + "summary": ( + "The four-person Artemis II crew has reached Earth orbit " + "after a successful launch from Kennedy Space Center." + ), + }] + + new = [{ + "id": "new-artemis-flyby", + "title": "Artemis II crew completes historic lunar flyby", + "summary": ( + "The Artemis II crew completed their lunar flyby today, " + "becoming the first humans to orbit the Moon since Apollo 17. " + "All systems nominal for return trajectory." + ), + }] + + duplicates = dedup.find_duplicates(new, recent) + assert "new-artemis-flyby" not in duplicates, ( + "Different mission milestones should NOT be flagged as duplicate" + ) + + +class TestBatchBehavior: + """Test that batched comparison works correctly with mixed results.""" + + def test_mixed_batch_duplicates_and_updates(self, dedup): + """Batch with both duplicates and updates — only duplicates filtered.""" + recent = [ + { + "id": "recent-soleimani", + "title": "US revokes residency of Soleimani relatives, detains two", + "summary": "The US revoked green cards for Iranian nationals tied to Iran's leadership.", + }, + { + "id": "recent-hormuz", + "title": "President Trump gives Iran 48-hour Hormuz deadline", + "summary": "Trump said Iran had 48 hours to reopen the Strait of Hormuz.", + }, + ] + + new = [ + { + "id": "new-soleimani-dupe", + "title": "US revokes visas of Soleimani relatives, detains two", + "summary": "Trump administration revoked green cards of Soleimani relatives. ICE detained two in California.", + }, + { + "id": "new-hormuz-update", + "title": "Iran allows Iraqi ships through Strait of Hormuz", + "summary": "Iran announced Iraqi vessels can transit Hormuz despite broader restrictions.", + }, + ] + + duplicates = dedup.find_duplicates(new, recent) + + assert "new-soleimani-dupe" in duplicates, ( + "Soleimani rewrite should be caught as duplicate" + ) + assert "new-hormuz-update" not in duplicates, ( + "Hormuz development should pass through as a new story" + ) diff --git a/aggregators/kagi-news/tests/test_state_manager.py b/aggregators/kagi-news/tests/test_state_manager.py index 3723041..f6d275e 100644 --- a/aggregators/kagi-news/tests/test_state_manager.py +++ b/aggregators/kagi-news/tests/test_state_manager.py @@ -225,3 +225,104 @@ class TestStateManager: # Old entry should be gone assert not manager.is_posted(feed_url, "old-guid") assert manager.is_posted(feed_url, "new-guid") + + def test_mark_posted_stores_title_and_summary(self, temp_state_file): + """Test that mark_posted stores title and summary_snippet.""" + manager = StateManager(temp_state_file) + feed_url = "https://news.kagi.com/world.xml" + + manager.mark_posted( + feed_url, "guid-1", "at://test/1", + title="US announces tariffs", + summary_snippet="The United States has announced new tariffs on goods." + ) + + state_data = json.loads(temp_state_file.read_text()) + entry = state_data['feeds'][feed_url]['posted_guids'][0] + assert entry['title'] == "US announces tariffs" + assert entry['summary_snippet'] == "The United States has announced new tariffs on goods." + + def test_mark_posted_truncates_summary_to_200(self, temp_state_file): + """Test that summary_snippet is truncated to 200 characters.""" + manager = StateManager(temp_state_file) + feed_url = "https://news.kagi.com/world.xml" + + long_summary = "x" * 300 + manager.mark_posted( + feed_url, "guid-1", "at://test/1", + title="Title", summary_snippet=long_summary + ) + + state_data = json.loads(temp_state_file.read_text()) + entry = state_data['feeds'][feed_url]['posted_guids'][0] + assert len(entry['summary_snippet']) == 200 + + def test_get_recent_stories_within_window(self, temp_state_file): + """Test get_recent_stories returns entries within lookback window.""" + manager = StateManager(temp_state_file) + feed_url = "https://news.kagi.com/world.xml" + + # Add a recent story with title + manager.mark_posted( + feed_url, "guid-1", "at://test/1", + title="Recent Story", summary_snippet="Recent summary" + ) + + recent = manager.get_recent_stories(feed_url, days=4) + assert len(recent) == 1 + assert recent[0]['id'] == "guid-1" + assert recent[0]['title'] == "Recent Story" + assert recent[0]['summary'] == "Recent summary" + + def test_get_recent_stories_excludes_old_entries(self, temp_state_file): + """Test that get_recent_stories excludes entries outside lookback.""" + manager = StateManager(temp_state_file) + feed_url = "https://news.kagi.com/world.xml" + + # Manually add an old entry with title + old_timestamp = (datetime.now() - timedelta(days=5)).isoformat() + state_data = { + 'feeds': { + feed_url: { + 'posted_guids': [{ + 'guid': 'old-guid', + 'post_uri': 'at://test/old', + 'posted_at': old_timestamp, + 'title': 'Old Story', + 'summary_snippet': 'Old summary' + }], + 'last_successful_run': None + } + } + } + temp_state_file.write_text(json.dumps(state_data, indent=2)) + + manager = StateManager(temp_state_file) + recent = manager.get_recent_stories(feed_url, days=4) + assert len(recent) == 0 + + def test_get_recent_stories_skips_entries_without_title(self, temp_state_file): + """Test backward compat: entries without title are skipped.""" + manager = StateManager(temp_state_file) + feed_url = "https://news.kagi.com/world.xml" + + # Old-format entry (no title) + manager.mark_posted(feed_url, "old-format-guid", "at://test/1") + + # New-format entry (with title) + manager.mark_posted( + feed_url, "new-format-guid", "at://test/2", + title="Story Title", summary_snippet="Summary" + ) + + recent = manager.get_recent_stories(feed_url, days=4) + assert len(recent) == 1 + assert recent[0]['id'] == "new-format-guid" + + def test_get_recent_stories_empty_feed(self, temp_state_file): + """Test get_recent_stories with no entries.""" + manager = StateManager(temp_state_file) + feed_url = "https://news.kagi.com/world.xml" + + recent = manager.get_recent_stories(feed_url, days=4) + assert recent == [] -- 2.51.2