diff --git a/aggregators/reddit-highlights/.gitignore b/aggregators/reddit-highlights/.gitignore new file mode 100644 index 0000000..783188a --- /dev/null +++ b/aggregators/reddit-highlights/.gitignore @@ -0,0 +1,41 @@ +# Environment and config +.env +config.yaml +venv/ + +# State files +data/*.json +data/world.xml + +# Python +__pycache__/ +*.py[cod] +*$py.class +*.so +.Python +build/ +develop-eggs/ +dist/ +downloads/ +eggs/ +.eggs/ +lib/ +lib64/ +parts/ +sdist/ +var/ +wheels/ +*.egg-info/ +.installed.cfg +*.egg + +# Testing +.pytest_cache/ +.coverage +htmlcov/ + +# IDE +.vscode/ +.idea/ +*.swp +*.swo diff --git a/aggregators/reddit-highlights/README.md b/aggregators/reddit-highlights/README.md new file mode 100644 index 0000000..ce5aa27 --- /dev/null +++ b/aggregators/reddit-highlights/README.md @@ -0,0 +1,140 @@ +# Reddit Highlights Aggregator + +Aggregates video highlights from Reddit subreddits (e.g., r/nba) and posts them to Coves communities. + +## Features + +- Fetches posts from Reddit via RSS (no API key required) +- Extracts streamable.com video links +- Posts to configured Coves communities with proper attribution +- Anti-detection jitter (randomized polling intervals) +- State tracking for deduplication +- Docker deployment with cron scheduler + +## Quick Start + +1. **Copy environment file:** + ```bash + cp .env.example .env + ``` + +2. **Configure your Coves API key:** + ```bash + # Edit .env and set your API key + COVES_API_KEY=ckapi_your_key_here + ``` + +3. **Build and run:** + ```bash + docker-compose up -d + ``` + +4. **View logs:** + ```bash + docker-compose logs -f + ``` + +## Configuration + +### config.yaml + +```yaml +coves_api_url: "https://coves.social" + +subreddits: + - name: "nba" + community_handle: "nba.coves.social" + enabled: true + +allowed_domains: + - streamable.com +``` + +### Adding More Subreddits + +1. Add entry to `config.yaml`: + ```yaml + - name: "soccer" + community_handle: "soccer.coves.social" + enabled: true + ``` + +2. Authorize the aggregator for the new community in Coves + +3. Restart the container: + ```bash + docker-compose restart + ``` + +## Polling Schedule + +- Cron runs every **10 minutes** +- Python script adds **0-10 minutes random jitter** +- Effective polling interval: **10-20 minutes** (varies each run) + +This randomization helps avoid bot detection patterns. + +## Development + +### Setup + +```bash +# Create virtual environment +python -m venv venv +source venv/bin/activate + +# Install dependencies +pip install -r requirements.txt +``` + +### Run Tests + +```bash +pytest +``` + +### Run Manually + +```bash +# Set environment variables +export COVES_API_KEY=ckapi_your_key +export SKIP_JITTER=true # Skip delay for testing + +# Run +python -m src.main +``` + +## Architecture + +``` +src/ +├── main.py # Orchestration (CRON entry point) +├── rss_fetcher.py # RSS feed fetching with retry +├── link_extractor.py # Streamable URL detection +├── coves_client.py # Coves API client +├── state_manager.py # Deduplication state tracking +├── config.py # YAML config loader +└── models.py # Data models +``` + +## Post Format + +Posts are created with: +- **Title**: Reddit post title +- **Embed**: Streamable video link with metadata +- **Sources**: Link back to original Reddit post + +Example embed: +```json +{ + "$type": "social.coves.embed.external", + "external": { + "uri": "https://streamable.com/abc123", + "title": "LeBron with the chase-down block!", + "description": "From r/nba", + "sources": [ + {"uri": "https://reddit.com/...", "title": "r/nba", "domain": "reddit.com"} + ] + } +} +``` diff --git a/aggregators/reddit-highlights/config.example.yaml b/aggregators/reddit-highlights/config.example.yaml new file mode 100644 index 0000000..ab2b9ff --- /dev/null +++ b/aggregators/reddit-highlights/config.example.yaml @@ -0,0 +1,39 @@ +# Reddit Highlights Aggregator Configuration +# +# Copy this file to config.yaml and customize for your setup. +# +# This file configures which subreddits to fetch video highlights from +# and which Coves communities to post them to. + +# Coves API endpoint (can be overridden with COVES_API_URL env var) +coves_api_url: "https://coves.social" + +# Subreddit-to-community mappings +# Add entries here for each subreddit you want to aggregate +subreddits: + # NBA highlights + - name: "nba" + community_handle: "nba.coves.social" + enabled: true + + # Example: Soccer/football highlights (disabled by default) + # - name: "soccer" + # community_handle: "soccer.coves.social" + # enabled: false + + # Example: NHL highlights (disabled by default) + # - name: "hockey" + # community_handle: "hockey.coves.social" + # enabled: false + +# Allowed video hosting domains +# Only posts with links to these domains will be imported +# Add more domains here as needed +allowed_domains: + - streamable.com + # Future options: + # - streamff.com + # - streamja.com + +# Logging level (debug, info, warning, error) +log_level: "info" diff --git a/aggregators/reddit-highlights/pytest.ini b/aggregators/reddit-highlights/pytest.ini new file mode 100644 index 0000000..9855d94 --- /dev/null +++ b/aggregators/reddit-highlights/pytest.ini @@ -0,0 +1,6 @@ +[pytest] +testpaths = tests +python_files = test_*.py +python_classes = Test* +python_functions = test_* +addopts = -v --tb=short diff --git a/aggregators/reddit-highlights/requirements.txt b/aggregators/reddit-highlights/requirements.txt new file mode 100644 index 0000000..80d9d6d --- /dev/null +++ b/aggregators/reddit-highlights/requirements.txt @@ -0,0 +1,15 @@ +# Core dependencies +feedparser==6.0.11 +requests==2.31.0 +pyyaml==6.0.1 + +# Testing +pytest==8.1.1 +pytest-cov==5.0.0 +responses==0.25.0 + +# Development +black==24.3.0 +mypy==1.9.0 +types-PyYAML==6.0.12.12 +types-requests==2.31.0.20240311 diff --git a/aggregators/reddit-highlights/src/__init__.py b/aggregators/reddit-highlights/src/__init__.py new file mode 100644 index 0000000..a0a6292 --- /dev/null +++ b/aggregators/reddit-highlights/src/__init__.py @@ -0,0 +1 @@ +"""Reddit Highlights Aggregator for Coves.""" diff --git a/aggregators/reddit-highlights/src/config.py b/aggregators/reddit-highlights/src/config.py new file mode 100644 index 0000000..c723d28 --- /dev/null +++ b/aggregators/reddit-highlights/src/config.py @@ -0,0 +1,181 @@ +""" +Configuration Loader for Reddit Highlights Aggregator. + +Loads and validates configuration from YAML files. +""" +import os +import logging +from pathlib import Path +from typing import Dict, Any, List +import yaml +from urllib.parse import urlparse + +from src.models import AggregatorConfig, SubredditConfig, LogLevel + +logger = logging.getLogger(__name__) + + +class ConfigError(Exception): + """Configuration error.""" + + pass + + +class ConfigLoader: + """ + Loads and validates aggregator configuration. + + Supports: + - Loading from YAML file + - Environment variable overrides + - Validation of required fields + """ + + def __init__(self, config_path: Path): + """ + Initialize config loader. + + Args: + config_path: Path to config.yaml file + """ + self.config_path = Path(config_path) + + def load(self) -> AggregatorConfig: + """ + Load and validate configuration. + + Returns: + AggregatorConfig object + + Raises: + ConfigError: If config is invalid or missing + """ + if not self.config_path.exists(): + raise ConfigError(f"Configuration file not found: {self.config_path}") + + try: + with open(self.config_path, "r") as f: + config_data = yaml.safe_load(f) + except yaml.YAMLError as e: + raise ConfigError(f"Failed to parse YAML: {e}") + + if not config_data: + raise ConfigError("Configuration file is empty") + + try: + return self._parse_config(config_data) + except ConfigError: + raise + except Exception as e: + raise ConfigError(f"Invalid configuration: {e}") + + def _parse_config(self, data: Dict[str, Any]) -> AggregatorConfig: + """ + Parse and validate configuration data. + + Args: + data: Parsed YAML data + + Returns: + AggregatorConfig object + + Raises: + ConfigError: If validation fails + """ + coves_api_url = os.getenv("COVES_API_URL", data.get("coves_api_url")) + if not coves_api_url: + raise ConfigError("Missing required field: coves_api_url") + + if not self._is_valid_url(coves_api_url): + raise ConfigError(f"Invalid URL for coves_api_url: {coves_api_url}") + + # Parse log level with validation + log_level_str = data.get("log_level", "info").lower() + try: + log_level = LogLevel(log_level_str) + except ValueError: + valid_levels = [level.value for level in LogLevel] + raise ConfigError(f"Invalid log_level '{log_level_str}'. Valid values: {valid_levels}") + + subreddits_data = data.get("subreddits", []) + if not subreddits_data: + raise ConfigError("Configuration must include at least one subreddit") + + subreddits = [] + for sub_data in subreddits_data: + subreddit = self._parse_subreddit(sub_data) + subreddits.append(subreddit) + + allowed_domains = tuple(data.get("allowed_domains", ["streamable.com"])) + + enabled_count = sum(1 for s in subreddits if s.enabled) + logger.info( + f"Loaded configuration with {len(subreddits)} subreddits ({enabled_count} enabled)" + ) + + return AggregatorConfig( + coves_api_url=coves_api_url, + subreddits=tuple(subreddits), # Convert to tuple for immutability + allowed_domains=allowed_domains, + log_level=log_level, + ) + + def _parse_subreddit(self, data: Dict[str, Any]) -> SubredditConfig: + """ + Parse and validate a single subreddit configuration. + + Args: + data: Subreddit configuration data + + Returns: + SubredditConfig object + + Raises: + ConfigError: If validation fails + """ + required_fields = ["name", "community_handle"] + for field in required_fields: + if field not in data: + raise ConfigError( + f"Missing required field in subreddit config: {field}" + ) + + name = data["name"] + community_handle = data["community_handle"] + enabled = data.get("enabled", True) + + if not name or not name.strip(): + raise ConfigError("Subreddit name cannot be empty") + + if not community_handle or not community_handle.strip(): + raise ConfigError(f"Community handle cannot be empty for subreddit '{name}'") + + return SubredditConfig( + name=name.strip().lower(), + community_handle=community_handle.strip(), + enabled=enabled, + ) + + def _is_valid_url(self, url: str) -> bool: + """ + Validate URL format. + + Only allows http and https schemes to prevent dangerous schemes + like file://, javascript://, or data:// URIs. + + Args: + url: URL to validate + + Returns: + True if valid HTTP/HTTPS URL, False otherwise + """ + try: + result = urlparse(url) + # Only allow http and https schemes + if result.scheme not in ("http", "https"): + logger.warning(f"URL has invalid scheme '{result.scheme}': {url}") + return False + return bool(result.netloc) + except ValueError as e: + logger.warning(f"Failed to parse URL '{url}': {e}") + return False diff --git a/aggregators/reddit-highlights/src/coves_client.py b/aggregators/reddit-highlights/src/coves_client.py new file mode 100644 index 0000000..b3d4620 --- /dev/null +++ b/aggregators/reddit-highlights/src/coves_client.py @@ -0,0 +1,285 @@ +""" +Coves API Client for posting to communities. + +Handles API key authentication and posting via XRPC. +""" +import logging +import requests +from typing import Dict, List, Optional + +logger = logging.getLogger(__name__) + + +class CovesAPIError(Exception): + """Base exception for Coves API errors.""" + + def __init__(self, message: str, status_code: int = None, response_body: str = None): + super().__init__(message) + self.status_code = status_code + self.response_body = response_body + + +class CovesAuthenticationError(CovesAPIError): + """Raised when authentication fails (401 Unauthorized).""" + pass + + +class CovesNotFoundError(CovesAPIError): + """Raised when a resource is not found (404 Not Found).""" + pass + + +class CovesRateLimitError(CovesAPIError): + """Raised when rate limit is exceeded (429 Too Many Requests).""" + pass + + +class CovesForbiddenError(CovesAPIError): + """Raised when access is forbidden (403 Forbidden).""" + pass + + +class CovesClient: + """ + Client for posting to Coves communities via XRPC. + + Handles: + - API key authentication + - Creating posts in communities (social.coves.community.post.create) + - External embed formatting + """ + + # API key format constants (must match Go constants in apikey_service.go) + API_KEY_PREFIX = "ckapi_" + API_KEY_TOTAL_LENGTH = 70 # 6 (prefix) + 64 (32 bytes hex-encoded) + + def __init__(self, api_url: str, api_key: str): + """ + Initialize Coves client with API key authentication. + + Args: + api_url: Coves API URL for posting (e.g., "https://coves.social") + api_key: Coves API key, 70 characters total (6-char prefix + 64-char hex token) + + Raises: + ValueError: If api_key is empty, has wrong prefix, or wrong length + """ + # Validate API key format for early failure with clear error + if not api_key: + raise ValueError("API key cannot be empty") + if not api_key.startswith(self.API_KEY_PREFIX): + raise ValueError(f"API key must start with '{self.API_KEY_PREFIX}'") + if len(api_key) != self.API_KEY_TOTAL_LENGTH: + raise ValueError( + f"API key must be {self.API_KEY_TOTAL_LENGTH} characters " + f"(got {len(api_key)})" + ) + + self.api_url = api_url.rstrip('/') + self.api_key = api_key + self.session = requests.Session() + self.session.headers['Authorization'] = f'Bearer {api_key}' + self.session.headers['Content-Type'] = 'application/json' + + def authenticate(self): + """ + No-op for API key authentication. + + API key is set in the session headers during initialization. + This method is kept for backward compatibility with existing code + that calls authenticate() before making requests. + """ + logger.info("Using API key authentication (no session creation needed)") + + def create_post( + self, + community_handle: str, + content: str, + facets: List[Dict], + title: Optional[str] = None, + embed: Optional[Dict] = None, + thumbnail_url: Optional[str] = None + ) -> str: + """ + Create a post in a community. + + Args: + community_handle: Community handle (e.g., "world-news.coves.social") + content: Post content (rich text) + facets: Rich text facets (formatting, links) + title: Optional post title + embed: Optional external embed + thumbnail_url: Optional thumbnail URL (for trusted aggregators only) + + Returns: + AT Proto URI of created post (e.g., "at://did:plc:.../social.coves.post/...") + + Raises: + CovesAuthenticationError: If authentication fails (401) + CovesForbiddenError: If access is denied (403) + CovesNotFoundError: If community not found (404) + CovesRateLimitError: If rate limit exceeded (429) + CovesAPIError: For other API errors or invalid responses + requests.RequestException: For network-level errors + """ + try: + # Prepare post data for social.coves.community.post.create endpoint + post_data = { + "community": community_handle, + "content": content, + "facets": facets + } + + # Add title if provided + if title: + post_data["title"] = title + + # Add embed if provided + if embed: + post_data["embed"] = embed + + # Add thumbnail URL at top level if provided (for trusted aggregators) + if thumbnail_url: + post_data["thumbnailUrl"] = thumbnail_url + + # Use Coves-specific endpoint (not direct PDS write) + # This provides validation, authorization, and business logic + logger.info(f"Creating post in community: {community_handle}") + + # Make HTTP request to XRPC endpoint using session with API key + url = f"{self.api_url}/xrpc/social.coves.community.post.create" + response = self.session.post(url, json=post_data, timeout=30) + + # Handle specific error cases + if not response.ok: + # Log status code but not full response body (may contain sensitive data) + logger.error(f"Post creation failed with status {response.status_code}") + self._raise_for_status(response) + + try: + result = response.json() + post_uri = result["uri"] + except (ValueError, KeyError) as e: + # ValueError for invalid JSON, KeyError for missing 'uri' field + logger.error(f"Failed to parse post creation response: {e}") + raise CovesAPIError( + f"Invalid response from server: {e}", + status_code=response.status_code, + response_body=response.text + ) + + logger.info(f"Post created: {post_uri}") + return post_uri + + except requests.RequestException as e: + # Network errors, timeouts, etc. + logger.error(f"Network error creating post: {e}") + raise + except CovesAPIError: + # Re-raise our custom exceptions as-is + raise + + def create_external_embed( + self, + uri: str, + title: str, + description: str, + sources: Optional[List[Dict]] = None, + embed_type: Optional[str] = None, + provider: Optional[str] = None, + domain: Optional[str] = None + ) -> Dict: + """ + Create external embed object for hot-linked content. + + Args: + uri: URL of the external content + title: Title of the content + description: Description/summary + sources: Optional list of source dicts with uri, title, domain + embed_type: Type hint for rendering (article, image, video, website) + provider: Service provider name (e.g., streamable, imgur) + domain: Domain of the linked content (e.g., streamable.com) + + Returns: + Embed dictionary ready for post creation + """ + external = { + "uri": uri, + "title": title, + "description": description + } + + if sources: + external["sources"] = sources + + if embed_type: + external["embedType"] = embed_type + + if provider: + external["provider"] = provider + + if domain: + external["domain"] = domain + + return { + "$type": "social.coves.embed.external", + "external": external + } + + def _raise_for_status(self, response: requests.Response) -> None: + """ + Raise specific exceptions based on HTTP status code. + + Args: + response: The HTTP response object + + Raises: + CovesAuthenticationError: For 401 Unauthorized + CovesNotFoundError: For 404 Not Found + CovesRateLimitError: For 429 Too Many Requests + CovesAPIError: For other 4xx/5xx errors + """ + status_code = response.status_code + error_body = response.text + + if status_code == 401: + raise CovesAuthenticationError( + f"Authentication failed: {error_body}", + status_code=status_code, + response_body=error_body + ) + elif status_code == 403: + raise CovesForbiddenError( + f"Access forbidden: {error_body}", + status_code=status_code, + response_body=error_body + ) + elif status_code == 404: + raise CovesNotFoundError( + f"Resource not found: {error_body}", + status_code=status_code, + response_body=error_body + ) + elif status_code == 429: + raise CovesRateLimitError( + f"Rate limit exceeded: {error_body}", + status_code=status_code, + response_body=error_body + ) + else: + raise CovesAPIError( + f"API request failed ({status_code}): {error_body}", + status_code=status_code, + response_body=error_body + ) + + def _get_timestamp(self) -> str: + """ + Get current timestamp in ISO 8601 format. + + Returns: + ISO timestamp string + """ + from datetime import datetime, timezone + return datetime.now(timezone.utc).isoformat().replace("+00:00", "Z") diff --git a/aggregators/reddit-highlights/src/link_extractor.py b/aggregators/reddit-highlights/src/link_extractor.py new file mode 100644 index 0000000..3615838 --- /dev/null +++ b/aggregators/reddit-highlights/src/link_extractor.py @@ -0,0 +1,212 @@ +""" +Link extractor for detecting streamable.com URLs from Reddit RSS entries. +""" +import re +import logging +from typing import List, Optional +from urllib.parse import urlparse + +logger = logging.getLogger(__name__) + + +class LinkExtractor: + """ + Extracts video links from Reddit RSS entries. + + Focuses on streamable.com but can be extended to other video hosts. + """ + + # Supported video hosting domains + DEFAULT_ALLOWED_DOMAINS = ["streamable.com"] + + # Pattern to find URLs in text/HTML content. + # Matches http:// or https:// followed by any non-whitespace characters + # that aren't typically URL delimiters in HTML/text (< > " ' ) ]). + # This is intentionally permissive to catch URLs in various contexts, + # with further validation done by is_allowed_url(). + URL_PATTERN = re.compile( + r'https?://[^\s<>"\')\]]+', + re.IGNORECASE, + ) + + def __init__(self, allowed_domains: Optional[List[str]] = None): + """ + Initialize link extractor. + + Args: + allowed_domains: List of allowed video hosting domains. + Defaults to ["streamable.com"] + """ + self.allowed_domains = allowed_domains or self.DEFAULT_ALLOWED_DOMAINS + # Normalize domains to lowercase + self.allowed_domains = [d.lower() for d in self.allowed_domains] + + def extract_video_url(self, entry) -> Optional[str]: + """ + Extract video URL from a Reddit RSS entry. + + Checks: + 1. Direct link (entry.link is a video URL) + 2. Entry content/description for embedded video URLs + + Args: + entry: feedparser entry object from Reddit RSS + + Returns: + Video URL if found, None otherwise + """ + # Check 1: Direct link + if hasattr(entry, "link") and entry.link: + if self.is_allowed_url(entry.link): + logger.debug(f"Found video URL in direct link: {entry.link}") + return self._normalize_url(entry.link) + + # Check 2: Entry content (Reddit RSS uses 'content' field) + if hasattr(entry, "content") and entry.content: + for content_item in entry.content: + if hasattr(content_item, "value"): + url = self._find_url_in_text(content_item.value) + if url: + logger.debug(f"Found video URL in content: {url}") + return url + + # Check 3: Entry description/summary + if hasattr(entry, "description") and entry.description: + url = self._find_url_in_text(entry.description) + if url: + logger.debug(f"Found video URL in description: {url}") + return url + + if hasattr(entry, "summary") and entry.summary: + url = self._find_url_in_text(entry.summary) + if url: + logger.debug(f"Found video URL in summary: {url}") + return url + + return None + + def _find_url_in_text(self, text: str) -> Optional[str]: + """ + Find first allowed video URL in text. + + Args: + text: Text/HTML content to search + + Returns: + First video URL found, or None + """ + if not text: + return None + + urls = self.URL_PATTERN.findall(text) + for url in urls: + # Clean up common trailing characters + url = url.rstrip(".,;:!?") + if self.is_allowed_url(url): + return self._normalize_url(url) + + return None + + def is_allowed_url(self, url: str) -> bool: + """ + Check if URL is from an allowed video hosting domain. + + Args: + url: URL to check + + Returns: + True if URL is from allowed domain + """ + if not url: + return False + + try: + parsed = urlparse(url) + domain = parsed.netloc.lower() + + # Remove www. prefix for comparison + if domain.startswith("www."): + domain = domain[4:] + + return domain in self.allowed_domains + + except ValueError as e: + logger.debug(f"Failed to parse URL '{url}': {e}") + return False + + def _normalize_url(self, url: str) -> str: + """ + Normalize URL for consistent storage/comparison. + + Args: + url: URL to normalize + + Returns: + Normalized URL + """ + # Remove trailing slashes + url = url.rstrip("/") + + # Ensure https + if url.startswith("http://"): + url = "https://" + url[7:] + + return url + + def get_video_id(self, url: str) -> Optional[str]: + """ + Extract video ID from streamable URL for deduplication. + + Args: + url: Streamable URL + + Returns: + Video ID or None + """ + if not url: + return None + + try: + parsed = urlparse(url) + # Streamable URLs are like: https://streamable.com/abc123 + path = parsed.path.strip("/") + if path: + # Return first path segment as video ID + return path.split("/")[0] + except ValueError as e: + logger.debug(f"Failed to extract video ID from URL '{url}': {e}") + + return None + + def get_thumbnail_url(self, url: str, timeout: int = 10) -> Optional[str]: + """ + Fetch thumbnail URL from Streamable's oembed API. + + Args: + url: Streamable video URL + timeout: Request timeout in seconds + + Returns: + Thumbnail URL or None if fetch fails + """ + if not url or not self.is_allowed_url(url): + return None + + import requests + + oembed_url = f"https://api.streamable.com/oembed?url={url}" + + try: + response = requests.get(oembed_url, timeout=timeout) + response.raise_for_status() + data = response.json() + thumbnail = data.get("thumbnail_url") + if thumbnail: + logger.debug(f"Fetched thumbnail for {url}: {thumbnail[:50]}...") + return thumbnail + except requests.RequestException as e: + logger.warning(f"Failed to fetch oembed for {url}: {e}") + return None + except (ValueError, KeyError) as e: + logger.warning(f"Failed to parse oembed response for {url}: {e}") + return None diff --git a/aggregators/reddit-highlights/src/main.py b/aggregators/reddit-highlights/src/main.py new file mode 100644 index 0000000..34cf697 --- /dev/null +++ b/aggregators/reddit-highlights/src/main.py @@ -0,0 +1,364 @@ +""" +Main Orchestration Script for Reddit Highlights Aggregator. + +Coordinates all components to: +1. Apply anti-detection jitter delay +2. Fetch Reddit RSS feeds +3. Extract streamable video links +4. Deduplicate via state tracking +5. Post to Coves communities + +Designed to run via CRON (single execution, then exit). +""" +import os +import re +import sys +import time +import random +import logging +from pathlib import Path +from datetime import datetime +from typing import Optional + +from src.config import ConfigLoader +from src.rss_fetcher import RSSFetcher +from src.link_extractor import LinkExtractor +from src.state_manager import StateManager +from src.coves_client import CovesClient +from src.models import RedditPost + +# Setup logging +logging.basicConfig( + level=logging.INFO, format="%(asctime)s - %(name)s - %(levelname)s - %(message)s" +) +logger = logging.getLogger(__name__) + +# Reddit RSS URL template +REDDIT_RSS_URL = "https://www.reddit.com/r/{subreddit}/.rss" + +# Anti-detection jitter range (0-10 minutes in seconds) +JITTER_MIN_SECONDS = 0 +JITTER_MAX_SECONDS = 600 + + +class Aggregator: + """ + Main aggregator orchestration. + + Coordinates all components to fetch, filter, and post video highlights. + """ + + def __init__( + self, + config_path: Path, + state_file: Path, + coves_client: Optional[CovesClient] = None, + skip_jitter: bool = False, + ): + """ + Initialize aggregator. + + Args: + config_path: Path to config.yaml + state_file: Path to state.json + coves_client: Optional CovesClient (for testing) + skip_jitter: Skip anti-detection delay (for testing) + """ + self.skip_jitter = skip_jitter + + # Load configuration + logger.info("Loading configuration...") + config_loader = ConfigLoader(config_path) + self.config = config_loader.load() + + # Initialize components + logger.info("Initializing components...") + self.rss_fetcher = RSSFetcher() + self.link_extractor = LinkExtractor( + allowed_domains=self.config.allowed_domains + ) + self.state_manager = StateManager(state_file) + self.state_file = state_file + + # Initialize Coves client (or use provided one for testing) + if coves_client: + self.coves_client = coves_client + else: + api_key = os.getenv("COVES_API_KEY") + if not api_key: + raise ValueError("COVES_API_KEY environment variable required") + + self.coves_client = CovesClient( + api_url=self.config.coves_api_url, api_key=api_key + ) + + def run(self): + """ + Run aggregator: apply jitter, fetch, filter, post, and update state. + + This is the main entry point for CRON execution. + """ + logger.info("=" * 60) + logger.info("Starting Reddit Highlights Aggregator") + logger.info("=" * 60) + + # Anti-detection jitter: random delay before starting + if not self.skip_jitter: + jitter_seconds = random.uniform(JITTER_MIN_SECONDS, JITTER_MAX_SECONDS) + logger.info( + f"Applying anti-detection jitter: sleeping for {jitter_seconds:.1f} seconds " + f"({jitter_seconds/60:.1f} minutes)" + ) + time.sleep(jitter_seconds) + + # Get enabled subreddits only + enabled_subreddits = [s for s in self.config.subreddits if s.enabled] + logger.info(f"Processing {len(enabled_subreddits)} enabled subreddits") + + # Authenticate once at the start + try: + self.coves_client.authenticate() + except Exception as e: + logger.error(f"Failed to authenticate: {e}") + logger.error("Cannot continue without authentication") + raise RuntimeError("Authentication failed") from e + + # Process each subreddit + for subreddit_config in enabled_subreddits: + try: + self._process_subreddit(subreddit_config) + except KeyboardInterrupt: + # Re-raise interrupt signals - don't suppress user abort + logger.info("Received interrupt signal, stopping...") + raise + except Exception as e: + # Log error but continue with other subreddits + logger.error( + f"Error processing subreddit '{subreddit_config.name}': {e}", + exc_info=True, + ) + continue + + logger.info("=" * 60) + logger.info("Aggregator run completed") + logger.info("=" * 60) + + def _process_subreddit(self, subreddit_config): + """ + Process a single subreddit. + + Args: + subreddit_config: SubredditConfig object + """ + subreddit_name = subreddit_config.name + community_handle = subreddit_config.community_handle + + # Sanitize subreddit name to prevent URL injection + # Only allow alphanumeric, underscores, and hyphens + if not re.match(r'^[a-zA-Z0-9_-]+$', subreddit_name): + raise ValueError(f"Invalid subreddit name: {subreddit_name}") + + logger.info(f"Processing subreddit: r/{subreddit_name} -> {community_handle}") + + # Build RSS URL + rss_url = REDDIT_RSS_URL.format(subreddit=subreddit_name) + + # Fetch RSS feed + try: + feed = self.rss_fetcher.fetch_feed(rss_url) + except Exception as e: + logger.error(f"Failed to fetch feed for r/{subreddit_name}: {e}") + raise + + # Check for feed errors + if feed.bozo: + bozo_exception = getattr(feed, 'bozo_exception', None) + logger.warning( + f"Feed for r/{subreddit_name} has parsing issues (bozo flag set): {bozo_exception}" + ) + + # Process entries + new_posts = 0 + skipped_posts = 0 + no_video_count = 0 + + for entry in feed.entries: + try: + # Extract video URL + video_url = self.link_extractor.extract_video_url(entry) + if not video_url: + no_video_count += 1 + continue # Skip posts without video links + + # Get entry ID for deduplication + entry_id = self._get_entry_id(entry) + if self.state_manager.is_posted(subreddit_name, entry_id): + skipped_posts += 1 + logger.debug(f"Skipping already-posted entry: {entry_id}") + continue + + # Parse entry into RedditPost + reddit_post = self._parse_entry(entry, subreddit_name, video_url) + + # Create embed with sources and video metadata + # Note: Thumbnail is fetched by backend via unfurl service + embed = self.coves_client.create_external_embed( + uri=reddit_post.streamable_url, + title=reddit_post.title, + description=f"From r/{subreddit_name}", + sources=[ + { + "uri": reddit_post.reddit_url, + "title": f"r/{subreddit_name}", + "domain": "reddit.com", + } + ], + embed_type="video", + provider="streamable", + domain="streamable.com", + ) + + # Post to community + try: + post_uri = self.coves_client.create_post( + community_handle=community_handle, + title=reddit_post.title, + content="", # No additional content needed + facets=[], + embed=embed, + ) + + # Mark as posted (only if successful) + self.state_manager.mark_posted(subreddit_name, entry_id, post_uri) + new_posts += 1 + logger.info(f"Posted: {reddit_post.title[:50]}... -> {post_uri}") + + except Exception as e: + # Don't update state if posting failed + logger.error(f"Failed to post '{reddit_post.title}': {e}") + continue + + except Exception as e: + # Log error but continue with other entries + logger.error(f"Error processing entry: {e}", exc_info=True) + continue + + # Update last run timestamp + self.state_manager.update_last_run(subreddit_name, datetime.now()) + + logger.info( + f"r/{subreddit_name}: {new_posts} new posts, {skipped_posts} duplicates, " + f"{no_video_count} without video" + ) + + def _get_entry_id(self, entry) -> str: + """ + Get unique identifier for RSS entry. + + Args: + entry: feedparser entry + + Returns: + Unique ID string + """ + # Reddit RSS uses 'id' field with format like 't3_abc123' + if hasattr(entry, "id") and entry.id: + return entry.id + + # Fallback to link + if hasattr(entry, "link") and entry.link: + return entry.link + + # Last resort: title hash (using SHA-256 for security) + if hasattr(entry, "title") and entry.title: + import hashlib + + logger.warning(f"Using fallback hash for entry ID (no id or link found)") + return hashlib.sha256(entry.title.encode()).hexdigest() + + raise ValueError("Cannot determine entry ID") + + def _parse_entry(self, entry, subreddit: str, video_url: str) -> RedditPost: + """ + Parse RSS entry into RedditPost object. + + Args: + entry: feedparser entry + subreddit: Subreddit name + video_url: Extracted video URL + + Returns: + RedditPost object + """ + # Get entry ID + entry_id = self._get_entry_id(entry) + + # Get title + title = entry.title if hasattr(entry, "title") else "Untitled" + + # Get Reddit permalink + reddit_url = entry.link if hasattr(entry, "link") else "" + + # Get author (Reddit RSS uses 'author' field) + author = "" + if hasattr(entry, "author"): + author = entry.author + elif hasattr(entry, "author_detail") and hasattr(entry.author_detail, "name"): + author = entry.author_detail.name + + # Get published date + published = None + if hasattr(entry, "published_parsed") and entry.published_parsed: + try: + published = datetime(*entry.published_parsed[:6]) + except (TypeError, ValueError) as e: + logger.warning(f"Failed to parse published date for entry: {e}") + + return RedditPost( + id=entry_id, + title=title, + link=entry.link if hasattr(entry, "link") else "", + reddit_url=reddit_url, + subreddit=subreddit, + author=author, + published=published, + streamable_url=video_url, + ) + + +def main(): + """ + Main entry point for command-line execution. + + Usage: + python -m src.main + """ + # Get paths from environment or use defaults + config_path = Path(os.getenv("CONFIG_PATH", "config.yaml")) + state_file = Path(os.getenv("STATE_FILE", "data/state.json")) + + # Check for skip jitter flag (for testing) + skip_jitter = os.getenv("SKIP_JITTER", "").lower() in ("true", "1", "yes") + + # Validate config file exists + if not config_path.exists(): + logger.error(f"Configuration file not found: {config_path}") + logger.error("Please create config.yaml (see config.example.yaml)") + sys.exit(1) + + # Create aggregator and run + try: + aggregator = Aggregator( + config_path=config_path, + state_file=state_file, + skip_jitter=skip_jitter, + ) + aggregator.run() + sys.exit(0) + except Exception as e: + logger.error(f"Aggregator failed: {e}", exc_info=True) + sys.exit(1) + + +if __name__ == "__main__": + main() diff --git a/aggregators/reddit-highlights/src/models.py b/aggregators/reddit-highlights/src/models.py new file mode 100644 index 0000000..054216a --- /dev/null +++ b/aggregators/reddit-highlights/src/models.py @@ -0,0 +1,90 @@ +""" +Data models for Reddit Highlights Aggregator. +""" +from dataclasses import dataclass, field +from enum import Enum +from typing import List, Optional, Tuple +from datetime import datetime +import re + + +class LogLevel(Enum): + """Valid log levels for aggregator configuration.""" + DEBUG = "debug" + INFO = "info" + WARNING = "warning" + ERROR = "error" + CRITICAL = "critical" + + +@dataclass +class RedditPost: + """ + Represents a Reddit post with video content. + + Parsed from Reddit RSS feed entries. + """ + + id: str # Reddit post ID (e.g., "t3_1abc123" or just the rkey) + title: str # Post title + link: str # Direct link to content (may be streamable URL) + reddit_url: str # Permalink to Reddit post + subreddit: str # Subreddit name (without r/) + author: str # Reddit username + published: Optional[datetime] = None # Post publication time + streamable_url: Optional[str] = None # Extracted streamable URL (if found) + + def __post_init__(self): + """Validate required fields.""" + if not self.id: + raise ValueError("RedditPost.id cannot be empty") + if not self.title: + raise ValueError("RedditPost.title cannot be empty") + if not self.subreddit: + raise ValueError("RedditPost.subreddit cannot be empty") + + +@dataclass(frozen=True) +class SubredditConfig: + """ + Configuration for a single subreddit source. + + Maps a subreddit to a Coves community. + Immutable (frozen) to prevent accidental modification. + """ + + name: str # Subreddit name (e.g., "nba") + community_handle: str # Coves community (e.g., "nba.coves.social") + enabled: bool = True # Whether to fetch from this subreddit + + def __post_init__(self): + """Validate configuration fields.""" + if not self.name or not self.name.strip(): + raise ValueError("SubredditConfig.name cannot be empty") + if not self.community_handle or not self.community_handle.strip(): + raise ValueError("SubredditConfig.community_handle cannot be empty") + # Validate subreddit name format (alphanumeric, underscores, hyphens only) + if not re.match(r'^[a-zA-Z0-9_-]+$', self.name): + raise ValueError(f"Invalid subreddit name format: {self.name}") + + +@dataclass(frozen=True) +class AggregatorConfig: + """ + Full aggregator configuration. + + Loaded from config.yaml. + Immutable (frozen) to prevent accidental modification after loading. + """ + + coves_api_url: str + subreddits: Tuple[SubredditConfig, ...] # Use tuple for immutability + allowed_domains: Tuple[str, ...] = ("streamable.com",) # Default tuple + log_level: LogLevel = LogLevel.INFO + + def __post_init__(self): + """Validate configuration.""" + if not self.coves_api_url: + raise ValueError("AggregatorConfig.coves_api_url cannot be empty") + if not self.subreddits: + raise ValueError("AggregatorConfig.subreddits cannot be empty") diff --git a/aggregators/reddit-highlights/src/rss_fetcher.py b/aggregators/reddit-highlights/src/rss_fetcher.py new file mode 100644 index 0000000..56384f6 --- /dev/null +++ b/aggregators/reddit-highlights/src/rss_fetcher.py @@ -0,0 +1,84 @@ +""" +RSS feed fetcher with retry logic and error handling. +""" +import time +import logging +import requests +import feedparser +from typing import Optional + +logger = logging.getLogger(__name__) + + +class RSSFetcher: + """ + Fetches and parses RSS feeds with retry logic and error handling. + + Features: + - Configurable timeout and retry count + - Exponential backoff on failures + - Custom User-Agent header (required by Reddit) + - Automatic HTTP to HTTPS upgrade handling + """ + + DEFAULT_USER_AGENT = "Coves-Reddit-Aggregator/1.0 (https://coves.social; contact@coves.social)" + + def __init__(self, timeout: int = 30, max_retries: int = 3, user_agent: Optional[str] = None): + """ + Initialize RSS fetcher. + + Args: + timeout: Request timeout in seconds + max_retries: Maximum number of retry attempts + user_agent: Custom User-Agent string (Reddit requires this) + """ + self.timeout = timeout + self.max_retries = max_retries + self.user_agent = user_agent or self.DEFAULT_USER_AGENT + + def fetch_feed(self, url: str) -> feedparser.FeedParserDict: + """ + Fetch and parse an RSS feed. + + Args: + url: RSS feed URL + + Returns: + Parsed feed object + + Raises: + ValueError: If URL is empty + requests.RequestException: If all retry attempts fail + """ + if not url: + raise ValueError("URL cannot be empty") + + last_error = None + + for attempt in range(self.max_retries): + try: + logger.info(f"Fetching feed from {url} (attempt {attempt + 1}/{self.max_retries})") + + headers = {"User-Agent": self.user_agent} + response = requests.get(url, timeout=self.timeout, headers=headers) + response.raise_for_status() + + # Parse with feedparser + feed = feedparser.parse(response.content) + + logger.info(f"Successfully fetched feed: {feed.feed.get('title', 'Unknown')}") + return feed + + except requests.RequestException as e: + last_error = e + logger.warning(f"Fetch attempt {attempt + 1} failed: {e}") + + if attempt < self.max_retries - 1: + # Exponential backoff + sleep_time = 2 ** attempt + logger.info(f"Retrying in {sleep_time} seconds...") + time.sleep(sleep_time) + + # All retries exhausted + logger.error(f"Failed to fetch feed after {self.max_retries} attempts") + raise last_error diff --git a/aggregators/reddit-highlights/src/state_manager.py b/aggregators/reddit-highlights/src/state_manager.py new file mode 100644 index 0000000..37cc7e8 --- /dev/null +++ b/aggregators/reddit-highlights/src/state_manager.py @@ -0,0 +1,255 @@ +""" +State Manager for tracking posted stories. + +Handles deduplication by tracking which stories have already been posted. +Uses JSON file for persistence. +""" +import json +import logging +from pathlib import Path +from datetime import datetime, timedelta +from typing import Optional, Dict, List + +logger = logging.getLogger(__name__) + + +class StateManager: + """ + Manages aggregator state for deduplication. + + Tracks posted Reddit entries per subreddit to prevent duplicate posting. + + Attributes tracked per subreddit: + - Posted GUIDs (with timestamps and Coves post URIs) + - Last successful run timestamp + - Automatic cleanup of old entries to prevent state file bloat + + Note: The 'feed_url' parameter in methods refers to the subreddit name + (e.g., 'nba'), not a full RSS URL. This naming is historical but the + functionality uses subreddit names as keys. + """ + + def __init__(self, state_file: Path, max_guids_per_feed: int = 100, max_age_days: int = 30): + """ + Initialize state manager. + + Args: + state_file: Path to JSON state file + max_guids_per_feed: Maximum GUIDs to keep per feed (default: 100) + max_age_days: Maximum age in days for GUIDs (default: 30) + """ + self.state_file = Path(state_file) + self.max_guids_per_feed = max_guids_per_feed + self.max_age_days = max_age_days + self.state = self._load_state() + + def _load_state(self) -> Dict: + """Load state from file, or create new state if file doesn't exist.""" + if not self.state_file.exists(): + logger.info(f"Creating new state file at {self.state_file}") + state = {'feeds': {}} + self._save_state(state) + return state + + try: + with open(self.state_file, 'r') as f: + state = json.load(f) + logger.info(f"Loaded state from {self.state_file}") + return state + except json.JSONDecodeError as e: + # Backup corrupted file before overwriting + backup_path = self.state_file.with_suffix('.json.corrupted') + logger.error(f"State file corrupted: {e}. Backing up to {backup_path}") + try: + import shutil + shutil.copy2(self.state_file, backup_path) + logger.info(f"Corrupted state file backed up to {backup_path}") + except OSError as backup_error: + logger.warning(f"Failed to backup corrupted state file: {backup_error}") + state = {'feeds': {}} + self._save_state(state) + return state + + def _save_state(self, state: Optional[Dict] = None): + """ + Save state to file atomically. + + Uses write-to-temp-then-rename pattern to prevent corruption + if the process is interrupted during write. + + Raises: + OSError: If write fails (after logging the error) + """ + if state is None: + state = self.state + + # Ensure parent directory exists + self.state_file.parent.mkdir(parents=True, exist_ok=True) + + # Write to temp file first for atomic update + temp_file = self.state_file.with_suffix('.json.tmp') + try: + with open(temp_file, 'w') as f: + json.dump(state, f, indent=2) + # Atomic rename (on POSIX systems) + temp_file.rename(self.state_file) + except OSError as e: + logger.error(f"Failed to save state file: {e}") + # Clean up temp file if it exists + if temp_file.exists(): + try: + temp_file.unlink() + except OSError: + pass + raise + + def _ensure_feed_exists(self, feed_url: str): + """Ensure feed entry exists in state.""" + if feed_url not in self.state['feeds']: + self.state['feeds'][feed_url] = { + 'posted_guids': [], + 'last_successful_run': None + } + + def is_posted(self, feed_url: str, guid: str) -> bool: + """ + Check if a story has already been posted. + + Args: + feed_url: RSS feed URL + guid: Story GUID + + Returns: + True if already posted, False otherwise + """ + self._ensure_feed_exists(feed_url) + + 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): + """ + Mark a story as posted. + + Args: + feed_url: RSS feed URL + guid: Story GUID + post_uri: AT Proto URI of created post + """ + self._ensure_feed_exists(feed_url) + + # Add to posted list + entry = { + 'guid': guid, + 'post_uri': post_uri, + 'posted_at': datetime.now().isoformat() + } + self.state['feeds'][feed_url]['posted_guids'].append(entry) + + # Auto-cleanup to keep state file manageable + self.cleanup_old_entries(feed_url) + + # Save state + self._save_state() + + logger.info(f"Marked as posted: {guid} -> {post_uri}") + + def get_last_run(self, feed_url: str) -> Optional[datetime]: + """ + Get last successful run timestamp for a feed. + + Args: + feed_url: RSS feed URL + + Returns: + Datetime of last run, or None if never run + """ + self._ensure_feed_exists(feed_url) + + timestamp_str = self.state['feeds'][feed_url]['last_successful_run'] + if timestamp_str is None: + return None + + return datetime.fromisoformat(timestamp_str) + + def update_last_run(self, feed_url: str, timestamp: datetime): + """ + Update last successful run timestamp. + + Args: + feed_url: RSS feed URL + timestamp: Timestamp of successful run + """ + self._ensure_feed_exists(feed_url) + + self.state['feeds'][feed_url]['last_successful_run'] = timestamp.isoformat() + self._save_state() + + logger.info(f"Updated last run for {feed_url}: {timestamp}") + + def cleanup_old_entries(self, feed_url: str): + """ + Remove old entries from state. + + Removes entries that are: + - Older than max_age_days + - Beyond max_guids_per_feed limit (keeps most recent) + + Args: + feed_url: RSS feed URL + """ + self._ensure_feed_exists(feed_url) + + posted_guids = self.state['feeds'][feed_url]['posted_guids'] + + # Filter out entries older than max_age_days + cutoff_date = datetime.now() - timedelta(days=self.max_age_days) + filtered = [] + for entry in posted_guids: + try: + posted_at = datetime.fromisoformat(entry['posted_at']) + if posted_at > cutoff_date: + filtered.append(entry) + except (KeyError, ValueError) as e: + # Skip entries with malformed or missing timestamps + logger.warning(f"Skipping entry with invalid timestamp: {e}") + continue + + # Keep only most recent max_guids_per_feed entries + # Sort by posted_at (most recent first) + filtered.sort(key=lambda x: x['posted_at'], reverse=True) + filtered = filtered[:self.max_guids_per_feed] + + # Update state + old_count = len(posted_guids) + new_count = len(filtered) + self.state['feeds'][feed_url]['posted_guids'] = filtered + + if old_count != new_count: + logger.info(f"Cleaned up {old_count - new_count} old entries for {feed_url}") + + def get_posted_count(self, feed_url: str) -> int: + """ + Get count of posted items for a feed. + + Args: + feed_url: RSS feed URL + + Returns: + Number of posted items + """ + self._ensure_feed_exists(feed_url) + return len(self.state['feeds'][feed_url]['posted_guids']) + + def get_all_posted_guids(self, feed_url: str) -> List[str]: + """ + Get all posted GUIDs for a feed. + + Args: + feed_url: RSS feed URL + + Returns: + List of GUIDs + """ + self._ensure_feed_exists(feed_url) + return [entry['guid'] for entry in self.state['feeds'][feed_url]['posted_guids']]