3803f3df4b
Setup wizard with consent-first cookie extraction (Chrome/Firefox/Safari), yt-dlp auto-install, ScrapeCreators push, quality scoring (5 core sources), status banner redesign, honest Reddit labeling, inline YouTube transcripts, Exa free web search, Reddit public fallback, and post-research quality nudge. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2126 lines
80 KiB
Python
2126 lines
80 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
last30days - Research a topic from the last 30 days on Reddit + X + YouTube + Web.
|
|
|
|
Usage:
|
|
python3 last30days.py <topic> [options]
|
|
|
|
Options:
|
|
--mock Use fixtures instead of real API calls
|
|
--emit=MODE Output mode: compact|json|md|context|path (default: compact)
|
|
--sources=MODE Source selection: auto|reddit|x|both (default: auto)
|
|
--quick Faster research with fewer sources (8-12 each)
|
|
--deep Comprehensive research with more sources (50-70 Reddit, 40-60 X)
|
|
--debug Enable verbose debug logging
|
|
--store Persist findings to SQLite database
|
|
--diagnose Show source availability diagnostics and exit
|
|
"""
|
|
|
|
import argparse
|
|
import atexit
|
|
import json
|
|
import os
|
|
import signal
|
|
import sys
|
|
import threading
|
|
from concurrent.futures import ThreadPoolExecutor, as_completed
|
|
from datetime import datetime, timezone
|
|
from pathlib import Path
|
|
|
|
# Add lib to path
|
|
SCRIPT_DIR = Path(__file__).parent.resolve()
|
|
sys.path.insert(0, str(SCRIPT_DIR))
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Global timeout & child process management
|
|
# ---------------------------------------------------------------------------
|
|
_child_pids: set = set()
|
|
_child_pids_lock = threading.Lock()
|
|
|
|
TIMEOUT_PROFILES = {
|
|
"quick": {"global": 90, "future": 30, "reddit_future": 60, "youtube_future": 60, "tiktok_future": 90, "instagram_future": 90, "hackernews_future": 30, "bluesky_future": 30, "truthsocial_future": 30, "polymarket_future": 15, "http": 15, "enrich_per": 8, "enrich_total": 30, "enrich_max_items": 10},
|
|
"default": {"global": 180, "future": 60, "reddit_future": 90, "youtube_future": 90, "tiktok_future": 120, "instagram_future": 120, "hackernews_future": 60, "bluesky_future": 60, "truthsocial_future": 60, "polymarket_future": 30, "http": 30, "enrich_per": 15, "enrich_total": 45, "enrich_max_items": 15},
|
|
"deep": {"global": 300, "future": 90, "reddit_future": 120, "youtube_future": 120, "tiktok_future": 150, "instagram_future": 150, "hackernews_future": 90, "bluesky_future": 90, "truthsocial_future": 90, "polymarket_future": 45, "http": 30, "enrich_per": 15, "enrich_total": 60, "enrich_max_items": 25},
|
|
}
|
|
|
|
# Valid source names for the --search flag
|
|
VALID_SEARCH_SOURCES = {
|
|
"reddit", "x", "hn", "bluesky", "bsky", "truthsocial", "truth", "youtube", "tiktok", "instagram",
|
|
"polymarket", "web", "xiaohongshu", "xhs",
|
|
}
|
|
|
|
|
|
def parse_search_flag(search_str: str) -> set:
|
|
"""Parse and validate the --search flag value.
|
|
|
|
Args:
|
|
search_str: Comma-separated source names (e.g. "reddit,hn")
|
|
|
|
Returns:
|
|
Set of validated source names
|
|
|
|
Raises:
|
|
SystemExit: If invalid sources are specified
|
|
"""
|
|
sources = set()
|
|
for s in search_str.split(","):
|
|
s = s.strip().lower()
|
|
if not s:
|
|
continue
|
|
if s == "xhs":
|
|
s = "xiaohongshu"
|
|
if s not in VALID_SEARCH_SOURCES:
|
|
print(
|
|
f"Error: Unknown search source '{s}'. "
|
|
f"Valid: {', '.join(sorted(VALID_SEARCH_SOURCES))}",
|
|
file=sys.stderr,
|
|
)
|
|
sys.exit(1)
|
|
sources.add(s)
|
|
if not sources:
|
|
print("Error: --search requires at least one source.", file=sys.stderr)
|
|
sys.exit(1)
|
|
return sources
|
|
|
|
|
|
def register_child_pid(pid: int):
|
|
"""Track a child process for cleanup."""
|
|
with _child_pids_lock:
|
|
_child_pids.add(pid)
|
|
|
|
|
|
def unregister_child_pid(pid: int):
|
|
"""Remove a child process from tracking."""
|
|
with _child_pids_lock:
|
|
_child_pids.discard(pid)
|
|
|
|
|
|
def _cleanup_children():
|
|
"""Kill all tracked child processes."""
|
|
with _child_pids_lock:
|
|
pids = list(_child_pids)
|
|
for pid in pids:
|
|
try:
|
|
os.killpg(os.getpgid(pid), signal.SIGTERM)
|
|
except (ProcessLookupError, PermissionError, OSError):
|
|
pass
|
|
|
|
|
|
atexit.register(_cleanup_children)
|
|
|
|
|
|
def _install_global_timeout(timeout_seconds: int):
|
|
"""Install a global timeout watchdog.
|
|
|
|
Uses SIGALRM on Unix, threading.Timer as fallback.
|
|
"""
|
|
if hasattr(signal, 'SIGALRM'):
|
|
def _handler(signum, frame):
|
|
sys.stderr.write(f"\n[TIMEOUT] Global timeout ({timeout_seconds}s) exceeded. Cleaning up.\n")
|
|
sys.stderr.flush()
|
|
_cleanup_children()
|
|
sys.exit(1)
|
|
signal.signal(signal.SIGALRM, _handler)
|
|
signal.alarm(timeout_seconds)
|
|
else:
|
|
# Windows fallback
|
|
def _watchdog():
|
|
sys.stderr.write(f"\n[TIMEOUT] Global timeout ({timeout_seconds}s) exceeded. Cleaning up.\n")
|
|
sys.stderr.flush()
|
|
_cleanup_children()
|
|
os._exit(1)
|
|
timer = threading.Timer(timeout_seconds, _watchdog)
|
|
timer.daemon = True
|
|
timer.start()
|
|
|
|
from lib import (
|
|
bird_x,
|
|
bluesky,
|
|
truthsocial,
|
|
dates,
|
|
dedupe,
|
|
hackernews,
|
|
xiaohongshu_api,
|
|
polymarket,
|
|
entity_extract,
|
|
env,
|
|
http,
|
|
models,
|
|
normalize,
|
|
openai_reddit,
|
|
reddit,
|
|
reddit_public,
|
|
reddit_enrich,
|
|
render,
|
|
schema,
|
|
score,
|
|
scrapecreators_x,
|
|
setup_wizard,
|
|
ui,
|
|
tiktok,
|
|
instagram,
|
|
websearch,
|
|
xai_x,
|
|
youtube_yt,
|
|
quality_nudge,
|
|
query_type as qt,
|
|
)
|
|
|
|
|
|
def load_fixture(name: str) -> dict:
|
|
"""Load a fixture file."""
|
|
fixture_path = SCRIPT_DIR.parent / "fixtures" / name
|
|
if fixture_path.exists():
|
|
with open(fixture_path) as f:
|
|
return json.load(f)
|
|
return {}
|
|
|
|
|
|
def _search_reddit(
|
|
topic: str,
|
|
config: dict,
|
|
selected_models: dict,
|
|
from_date: str,
|
|
to_date: str,
|
|
depth: str,
|
|
mock: bool,
|
|
) -> tuple:
|
|
"""Search Reddit (runs in thread).
|
|
|
|
Hierarchy:
|
|
1. ScrapeCreators (if SCRAPECREATORS_API_KEY exists) — premium, best quality
|
|
2. Public Reddit JSON (always available) — free, good for thread discovery
|
|
3. OpenAI Responses API — legacy fallback for backwards compatibility
|
|
|
|
Returns:
|
|
Tuple of (reddit_items, raw_response, error, used_scrapecreators)
|
|
"""
|
|
raw_response = None
|
|
reddit_error = None
|
|
used_scrapecreators = False
|
|
|
|
sc_token = config.get("SCRAPECREATORS_API_KEY")
|
|
|
|
if mock:
|
|
raw_response = load_fixture("openai_sample.json")
|
|
elif sc_token:
|
|
# === Tier 1: ScrapeCreators path (preferred) ===
|
|
used_scrapecreators = True
|
|
try:
|
|
sys.stderr.write("[Reddit] Using ScrapeCreators API\n")
|
|
sys.stderr.flush()
|
|
result = reddit.search_and_enrich(
|
|
topic, from_date, to_date,
|
|
depth=depth, token=sc_token,
|
|
)
|
|
reddit_items = result.get("items", [])
|
|
if result.get("error"):
|
|
reddit_error = result["error"]
|
|
return reddit_items, result, reddit_error, used_scrapecreators
|
|
except Exception as e:
|
|
reddit_error = f"ScrapeCreators: {type(e).__name__}: {e}"
|
|
sys.stderr.write(f"[Reddit] ScrapeCreators failed: {e}\n")
|
|
sys.stderr.flush()
|
|
used_scrapecreators = False
|
|
# Fall through to Tier 2 (public JSON)
|
|
|
|
# === Tier 2: Public Reddit JSON (free, always available) ===
|
|
if not mock:
|
|
try:
|
|
sys.stderr.write("[Reddit] Trying public Reddit JSON\n")
|
|
sys.stderr.flush()
|
|
reddit_items = reddit_public.search_reddit_public(
|
|
topic, from_date, to_date, depth=depth,
|
|
)
|
|
if reddit_items:
|
|
raw_response = {"source": "reddit_public", "items": reddit_items}
|
|
sys.stderr.write(f"[Reddit] Public JSON returned {len(reddit_items)} results\n")
|
|
sys.stderr.flush()
|
|
return reddit_items, raw_response, None, False
|
|
# Empty results — fall through to Tier 3
|
|
sys.stderr.write("[Reddit] Public JSON returned 0 results, trying OpenAI\n")
|
|
sys.stderr.flush()
|
|
except Exception as e:
|
|
sys.stderr.write(f"[Reddit] Public JSON failed: {e}\n")
|
|
sys.stderr.flush()
|
|
# Fall through to Tier 3
|
|
|
|
# === Tier 3: OpenAI Responses API (legacy fallback) ===
|
|
if not mock and config.get("OPENAI_API_KEY"):
|
|
try:
|
|
sys.stderr.write("[Reddit] Falling back to OpenAI Responses API\n")
|
|
sys.stderr.flush()
|
|
raw_response = openai_reddit.search_reddit(
|
|
config["OPENAI_API_KEY"],
|
|
selected_models["openai"],
|
|
topic,
|
|
from_date,
|
|
to_date,
|
|
depth=depth,
|
|
auth_source=config.get("OPENAI_AUTH_SOURCE", "api_key"),
|
|
account_id=config.get("OPENAI_CHATGPT_ACCOUNT_ID"),
|
|
)
|
|
except http.HTTPError as e:
|
|
raw_response = {"error": str(e)}
|
|
reddit_error = f"API error: {e}"
|
|
except Exception as e:
|
|
raw_response = {"error": str(e)}
|
|
reddit_error = f"{type(e).__name__}: {e}"
|
|
|
|
# Parse response (OpenAI path)
|
|
reddit_items = openai_reddit.parse_reddit_response(raw_response or {})
|
|
|
|
# Quick retry with simpler query if few results (OpenAI path only)
|
|
if len(reddit_items) < 5 and not mock and not reddit_error and config.get("OPENAI_API_KEY"):
|
|
core = openai_reddit._extract_core_subject(topic)
|
|
if core.lower() != topic.lower():
|
|
try:
|
|
retry_raw = openai_reddit.search_reddit(
|
|
config["OPENAI_API_KEY"],
|
|
selected_models["openai"],
|
|
core,
|
|
from_date, to_date,
|
|
depth=depth,
|
|
auth_source=config.get("OPENAI_AUTH_SOURCE", "api_key"),
|
|
account_id=config.get("OPENAI_CHATGPT_ACCOUNT_ID"),
|
|
)
|
|
retry_items = openai_reddit.parse_reddit_response(retry_raw)
|
|
existing_urls = {item.get("url") for item in reddit_items}
|
|
for item in retry_items:
|
|
if item.get("url") not in existing_urls:
|
|
reddit_items.append(item)
|
|
except Exception:
|
|
pass
|
|
|
|
# Subreddit-targeted fallback if still < 3 results (OpenAI path only)
|
|
if len(reddit_items) < 3 and not mock and not reddit_error and config.get("OPENAI_API_KEY"):
|
|
sub_query = openai_reddit._build_subreddit_query(topic)
|
|
try:
|
|
sub_raw = openai_reddit.search_reddit(
|
|
config["OPENAI_API_KEY"],
|
|
selected_models["openai"],
|
|
sub_query,
|
|
from_date, to_date,
|
|
depth=depth,
|
|
)
|
|
sub_items = openai_reddit.parse_reddit_response(sub_raw)
|
|
existing_urls = {item.get("url") for item in reddit_items}
|
|
for item in sub_items:
|
|
if item.get("url") not in existing_urls:
|
|
reddit_items.append(item)
|
|
except Exception:
|
|
pass
|
|
|
|
return reddit_items, raw_response, reddit_error, used_scrapecreators
|
|
|
|
|
|
def _search_x(
|
|
topic: str,
|
|
config: dict,
|
|
selected_models: dict,
|
|
from_date: str,
|
|
to_date: str,
|
|
depth: str,
|
|
mock: bool,
|
|
x_source: str = "xai",
|
|
) -> tuple:
|
|
"""Search X via Bird CLI or xAI (runs in thread).
|
|
|
|
Args:
|
|
x_source: 'bird' or 'xai' - which backend to use
|
|
|
|
Returns:
|
|
Tuple of (x_items, raw_response, error)
|
|
"""
|
|
raw_response = None
|
|
x_error = None
|
|
|
|
if mock:
|
|
raw_response = load_fixture("xai_sample.json")
|
|
x_items = xai_x.parse_x_response(raw_response or {})
|
|
return x_items, raw_response, x_error
|
|
|
|
# Use Bird if specified
|
|
if x_source == "bird":
|
|
try:
|
|
raw_response = bird_x.search_x(
|
|
topic,
|
|
from_date,
|
|
to_date,
|
|
depth=depth,
|
|
)
|
|
except Exception as e:
|
|
raw_response = {"error": str(e)}
|
|
x_error = f"{type(e).__name__}: {e}"
|
|
|
|
x_items = bird_x.parse_bird_response(raw_response or {}, query=topic)
|
|
|
|
# Check for error in response (Bird returns list on success, dict on error)
|
|
if raw_response and isinstance(raw_response, dict) and raw_response.get("error") and not x_error:
|
|
x_error = raw_response["error"]
|
|
|
|
return x_items, raw_response, x_error
|
|
|
|
# Use ScrapeCreators if specified
|
|
if x_source == "scrapecreators":
|
|
try:
|
|
raw_response = scrapecreators_x.search_x(
|
|
topic, from_date, to_date,
|
|
depth=depth,
|
|
token=config.get("SCRAPECREATORS_API_KEY"),
|
|
)
|
|
except Exception as e:
|
|
raw_response = {"error": str(e)}
|
|
x_error = f"{type(e).__name__}: {e}"
|
|
|
|
x_items = scrapecreators_x.parse_x_response(raw_response or {})
|
|
|
|
if raw_response and isinstance(raw_response, dict) and raw_response.get("error") and not x_error:
|
|
x_error = raw_response["error"]
|
|
|
|
return x_items, raw_response, x_error
|
|
|
|
# Use xAI (original behavior)
|
|
try:
|
|
raw_response = xai_x.search_x(
|
|
config["XAI_API_KEY"],
|
|
selected_models["xai"],
|
|
topic,
|
|
from_date,
|
|
to_date,
|
|
depth=depth,
|
|
)
|
|
except http.HTTPError as e:
|
|
raw_response = {"error": str(e)}
|
|
x_error = f"API error: {e}"
|
|
except Exception as e:
|
|
raw_response = {"error": str(e)}
|
|
x_error = f"{type(e).__name__}: {e}"
|
|
|
|
x_items = xai_x.parse_x_response(raw_response or {})
|
|
|
|
return x_items, raw_response, x_error
|
|
|
|
|
|
def _search_youtube(
|
|
topic: str,
|
|
from_date: str,
|
|
to_date: str,
|
|
depth: str,
|
|
) -> tuple:
|
|
"""Search YouTube via yt-dlp (runs in thread).
|
|
|
|
Returns:
|
|
Tuple of (youtube_items, youtube_error)
|
|
"""
|
|
youtube_error = None
|
|
|
|
try:
|
|
response = youtube_yt.search_and_transcribe(
|
|
topic, from_date, to_date, depth=depth,
|
|
)
|
|
except Exception as e:
|
|
return [], f"{type(e).__name__}: {e}"
|
|
|
|
youtube_items = youtube_yt.parse_youtube_response(response)
|
|
|
|
if response.get("error"):
|
|
youtube_error = response["error"]
|
|
|
|
return youtube_items, youtube_error
|
|
|
|
|
|
def _search_tiktok(
|
|
topic: str,
|
|
from_date: str,
|
|
to_date: str,
|
|
depth: str,
|
|
token: str,
|
|
) -> tuple:
|
|
"""Search TikTok via ScrapeCreators (runs in thread).
|
|
|
|
Returns:
|
|
Tuple of (tiktok_items, tiktok_error)
|
|
"""
|
|
tiktok_error = None
|
|
|
|
try:
|
|
response = tiktok.search_and_enrich(
|
|
topic, from_date, to_date, depth=depth, token=token,
|
|
)
|
|
except Exception as e:
|
|
return [], f"{type(e).__name__}: {e}"
|
|
|
|
tiktok_items = tiktok.parse_tiktok_response(response)
|
|
|
|
if response.get("error"):
|
|
tiktok_error = response["error"]
|
|
|
|
return tiktok_items, tiktok_error
|
|
|
|
|
|
def _search_instagram(
|
|
topic: str,
|
|
from_date: str,
|
|
to_date: str,
|
|
depth: str,
|
|
token: str,
|
|
) -> tuple:
|
|
"""Search Instagram via ScrapeCreators (runs in thread).
|
|
|
|
Returns:
|
|
Tuple of (instagram_items, instagram_error)
|
|
"""
|
|
instagram_error = None
|
|
|
|
try:
|
|
response = instagram.search_and_enrich(
|
|
topic, from_date, to_date, depth=depth, token=token,
|
|
)
|
|
except Exception as e:
|
|
return [], f"{type(e).__name__}: {e}"
|
|
|
|
instagram_items = instagram.parse_instagram_response(response)
|
|
|
|
if response.get("error"):
|
|
instagram_error = response["error"]
|
|
|
|
return instagram_items, instagram_error
|
|
|
|
|
|
def _search_hackernews(
|
|
topic: str,
|
|
from_date: str,
|
|
to_date: str,
|
|
depth: str,
|
|
) -> tuple:
|
|
"""Search Hacker News via Algolia (runs in thread).
|
|
|
|
Returns:
|
|
Tuple of (hn_items, hn_error)
|
|
"""
|
|
hn_error = None
|
|
|
|
try:
|
|
response = hackernews.search_hackernews(
|
|
topic, from_date, to_date, depth=depth,
|
|
)
|
|
except Exception as e:
|
|
return [], f"{type(e).__name__}: {e}"
|
|
|
|
hn_items = hackernews.parse_hackernews_response(response, query=topic)
|
|
|
|
if response.get("error"):
|
|
hn_error = response["error"]
|
|
|
|
return hn_items, hn_error
|
|
|
|
|
|
def _search_bluesky(
|
|
topic: str,
|
|
from_date: str,
|
|
to_date: str,
|
|
depth: str,
|
|
config: dict = None,
|
|
) -> tuple:
|
|
"""Search Bluesky via AT Protocol (runs in thread).
|
|
|
|
Returns:
|
|
Tuple of (bsky_items, bsky_error)
|
|
"""
|
|
bsky_error = None
|
|
|
|
try:
|
|
response = bluesky.search_bluesky(
|
|
topic, from_date, to_date, depth=depth, config=config,
|
|
)
|
|
except Exception as e:
|
|
return [], f"{type(e).__name__}: {e}"
|
|
|
|
bsky_items = bluesky.parse_bluesky_response(response)
|
|
|
|
if response.get("error"):
|
|
bsky_error = response["error"]
|
|
|
|
return bsky_items, bsky_error
|
|
|
|
|
|
def _search_truthsocial(
|
|
topic: str,
|
|
from_date: str,
|
|
to_date: str,
|
|
depth: str,
|
|
config: dict = None,
|
|
) -> tuple:
|
|
"""Search Truth Social via Mastodon API (runs in thread).
|
|
|
|
Returns:
|
|
Tuple of (ts_items, ts_error)
|
|
"""
|
|
ts_error = None
|
|
|
|
try:
|
|
response = truthsocial.search_truthsocial(
|
|
topic, from_date, to_date, depth=depth, config=config,
|
|
)
|
|
except Exception as e:
|
|
return [], f"{type(e).__name__}: {e}"
|
|
|
|
ts_items = truthsocial.parse_truthsocial_response(response)
|
|
|
|
if response.get("error"):
|
|
ts_error = response["error"]
|
|
|
|
return ts_items, ts_error
|
|
|
|
|
|
def _search_polymarket(
|
|
topic: str,
|
|
from_date: str,
|
|
to_date: str,
|
|
depth: str,
|
|
) -> tuple:
|
|
"""Search Polymarket via Gamma API (runs in thread).
|
|
|
|
Returns:
|
|
Tuple of (pm_items, pm_error)
|
|
"""
|
|
pm_error = None
|
|
|
|
try:
|
|
response = polymarket.search_polymarket(
|
|
topic, from_date, to_date, depth=depth,
|
|
)
|
|
except Exception as e:
|
|
return [], f"{type(e).__name__}: {e}"
|
|
|
|
pm_items = polymarket.parse_polymarket_response(response, topic=topic)
|
|
|
|
if response.get("error"):
|
|
pm_error = response["error"]
|
|
|
|
return pm_items, pm_error
|
|
|
|
|
|
def _search_web(
|
|
topic: str,
|
|
config: dict,
|
|
from_date: str,
|
|
to_date: str,
|
|
depth: str,
|
|
) -> tuple:
|
|
"""Search the web via native API backend (runs in thread).
|
|
|
|
Uses the best available backend: Parallel AI > Brave > OpenRouter.
|
|
|
|
Returns:
|
|
Tuple of (web_items, web_error)
|
|
web_items are raw dicts ready for websearch.normalize_websearch_items()
|
|
"""
|
|
from lib import brave_search, parallel_search, openrouter_search, exa_search
|
|
|
|
backend = env.get_web_search_source(config)
|
|
if not backend:
|
|
return [], "No web search API keys configured"
|
|
|
|
web_error = None
|
|
raw_results = []
|
|
|
|
try:
|
|
if backend == "exa":
|
|
raw_results = exa_search.search_web(
|
|
topic, from_date, to_date, config["EXA_API_KEY"], depth=depth,
|
|
)
|
|
elif backend == "parallel":
|
|
raw_results = parallel_search.search_web(
|
|
topic, from_date, to_date, config["PARALLEL_API_KEY"], depth=depth,
|
|
)
|
|
elif backend == "brave":
|
|
use_llm_ctx = os.environ.get("BRAVE_LLM_CONTEXT", "").strip() == "1"
|
|
raw_results = brave_search.search_web(
|
|
topic, from_date, to_date, config["BRAVE_API_KEY"],
|
|
depth=depth, use_llm_context=use_llm_ctx,
|
|
)
|
|
elif backend == "openrouter":
|
|
raw_results = openrouter_search.search_web(
|
|
topic, from_date, to_date, config["OPENROUTER_API_KEY"], depth=depth,
|
|
)
|
|
except Exception as e:
|
|
return [], f"{type(e).__name__}: {e}"
|
|
|
|
# Add IDs and date_confidence for websearch.normalize_websearch_items()
|
|
for i, item in enumerate(raw_results):
|
|
item.setdefault("id", f"W{i+1}")
|
|
if item.get("date") and not item.get("date_confidence"):
|
|
item["date_confidence"] = "med"
|
|
elif not item.get("date"):
|
|
item["date_confidence"] = "low"
|
|
item.setdefault("why_relevant", "")
|
|
|
|
return raw_results, web_error
|
|
|
|
|
|
def _search_xiaohongshu(
|
|
topic: str,
|
|
config: dict,
|
|
from_date: str,
|
|
to_date: str,
|
|
depth: str,
|
|
) -> tuple:
|
|
"""Search Xiaohongshu via xiaohongshu-mcp HTTP API (runs in thread).
|
|
|
|
Returns:
|
|
Tuple of (xiaohongshu_items, xiaohongshu_error)
|
|
Items are in web-item dict shape and can be normalized with websearch module.
|
|
"""
|
|
base_url = env.get_xiaohongshu_api_base(config)
|
|
try:
|
|
items = xiaohongshu_api.search_feeds(
|
|
topic=topic,
|
|
from_date=from_date,
|
|
to_date=to_date,
|
|
base_url=base_url,
|
|
depth=depth,
|
|
)
|
|
except Exception as e:
|
|
return [], f"{type(e).__name__}: {e}"
|
|
|
|
# Ensure all required keys exist for normalize_websearch_items()
|
|
for i, item in enumerate(items):
|
|
item.setdefault("id", f"XHS{i+1}")
|
|
item.setdefault("title", "")
|
|
item.setdefault("url", "")
|
|
item.setdefault("source_domain", "xiaohongshu.com")
|
|
item.setdefault("snippet", "")
|
|
if item.get("date") and not item.get("date_confidence"):
|
|
item["date_confidence"] = "med"
|
|
elif not item.get("date"):
|
|
item["date_confidence"] = "low"
|
|
item.setdefault("relevance", 0.5)
|
|
item.setdefault("why_relevant", "")
|
|
|
|
return items, None
|
|
|
|
|
|
def _run_supplemental(
|
|
topic: str,
|
|
reddit_items: list,
|
|
x_items: list,
|
|
from_date: str,
|
|
to_date: str,
|
|
depth: str,
|
|
x_source: str,
|
|
progress: ui.ProgressDisplay = None,
|
|
skip_reddit: bool = False,
|
|
resolved_handle: str = None,
|
|
) -> tuple:
|
|
"""Run Phase 2 supplemental searches based on entities from Phase 1.
|
|
|
|
Extracts handles/subreddits from initial results, then runs targeted
|
|
searches to find additional content the broad search missed.
|
|
|
|
Args:
|
|
topic: Original search topic
|
|
reddit_items: Phase 1 Reddit items (raw dicts)
|
|
x_items: Phase 1 X items (raw dicts)
|
|
from_date: Start date
|
|
to_date: End date
|
|
depth: Research depth
|
|
x_source: 'bird' or 'xai'
|
|
progress: Optional progress display
|
|
skip_reddit: If True, skip Reddit supplemental (e.g. rate-limited)
|
|
resolved_handle: X handle resolved by the agent (without @), searched unfiltered
|
|
|
|
Returns:
|
|
Tuple of (supplemental_reddit, supplemental_x)
|
|
"""
|
|
# Depth-dependent caps
|
|
if depth == "default":
|
|
max_handles = 3
|
|
max_subs = 3
|
|
count_per = 3
|
|
else: # deep
|
|
max_handles = 5
|
|
max_subs = 5
|
|
count_per = 5
|
|
|
|
# Extract entities from Phase 1 results
|
|
entities = entity_extract.extract_entities(
|
|
reddit_items, x_items,
|
|
max_handles=max_handles,
|
|
max_subreddits=max_subs,
|
|
)
|
|
|
|
has_handles = entities["x_handles"] and x_source == "bird"
|
|
has_subs = entities["reddit_subreddits"] and not skip_reddit
|
|
|
|
# Always run unfiltered search for resolved handle (even if entity-extracted).
|
|
# Entity-extracted handles get topic-filtered queries (from:handle topic),
|
|
# but resolved handles need UNFILTERED search (from:handle) to find posts
|
|
# that don't mention the topic string (e.g. Dor Brothers' viral tweet about
|
|
# Logan Paul doesn't contain "dor brothers" in the text).
|
|
has_resolved = bool(resolved_handle) and x_source == "bird"
|
|
|
|
if not has_handles and not has_subs and not has_resolved:
|
|
return [], []
|
|
|
|
parts = []
|
|
if has_resolved:
|
|
parts.append(f"@{resolved_handle} (resolved)")
|
|
if has_handles:
|
|
parts.append(f"@{', @'.join(entities['x_handles'][:3])}")
|
|
if has_subs:
|
|
parts.append(f"r/{', r/'.join(entities['reddit_subreddits'][:3])}")
|
|
sys.stderr.write(f"[Phase 2] Drilling into {' + '.join(parts)}\n")
|
|
sys.stderr.flush()
|
|
|
|
supplemental_reddit = []
|
|
supplemental_x = []
|
|
|
|
# Collect existing URLs to avoid adding duplicates before dedupe
|
|
existing_urls = set()
|
|
for item in reddit_items:
|
|
existing_urls.add(item.get("url", ""))
|
|
for item in x_items:
|
|
existing_urls.add(item.get("url", ""))
|
|
|
|
# Run supplemental searches in parallel
|
|
reddit_future = None
|
|
x_future = None
|
|
resolved_future = None
|
|
|
|
max_workers = sum([bool(has_subs), bool(has_handles), bool(has_resolved)])
|
|
with ThreadPoolExecutor(max_workers=max(max_workers, 1)) as executor:
|
|
if has_subs:
|
|
reddit_future = executor.submit(
|
|
openai_reddit.search_subreddits,
|
|
entities["reddit_subreddits"],
|
|
topic,
|
|
from_date,
|
|
to_date,
|
|
count_per,
|
|
)
|
|
|
|
if has_handles:
|
|
x_future = executor.submit(
|
|
bird_x.search_handles,
|
|
entities["x_handles"],
|
|
topic,
|
|
from_date,
|
|
count_per,
|
|
)
|
|
|
|
if has_resolved:
|
|
# Resolved handle: search unfiltered (topic=None) to get all recent posts
|
|
resolved_future = executor.submit(
|
|
bird_x.search_handles,
|
|
[resolved_handle],
|
|
None, # No topic filter - get all recent activity
|
|
from_date,
|
|
10, # More results for the topic entity
|
|
)
|
|
|
|
if reddit_future:
|
|
try:
|
|
raw_reddit = reddit_future.result(timeout=30)
|
|
# Filter out URLs already found in Phase 1
|
|
supplemental_reddit = [
|
|
item for item in raw_reddit
|
|
if item.get("url", "") not in existing_urls
|
|
]
|
|
except TimeoutError:
|
|
sys.stderr.write("[Phase 2] Supplemental Reddit timed out (30s)\n")
|
|
except Exception as e:
|
|
sys.stderr.write(f"[Phase 2] Supplemental Reddit error: {e}\n")
|
|
|
|
if x_future:
|
|
try:
|
|
raw_x = x_future.result(timeout=30)
|
|
supplemental_x = [
|
|
item for item in raw_x
|
|
if item.get("url", "") not in existing_urls
|
|
]
|
|
except TimeoutError:
|
|
sys.stderr.write("[Phase 2] Supplemental X timed out (30s)\n")
|
|
except Exception as e:
|
|
sys.stderr.write(f"[Phase 2] Supplemental X error: {e}\n")
|
|
|
|
if resolved_future:
|
|
try:
|
|
raw_resolved = resolved_future.result(timeout=30)
|
|
# Lower relevance for unfiltered handle posts (no topic keyword signal)
|
|
for item in raw_resolved:
|
|
item["relevance"] = 0.5
|
|
resolved_new = [
|
|
item for item in raw_resolved
|
|
if item.get("url", "") not in existing_urls
|
|
]
|
|
supplemental_x.extend(resolved_new)
|
|
if resolved_new:
|
|
sys.stderr.write(f"[Phase 2] +{len(resolved_new)} from @{resolved_handle}\n")
|
|
except TimeoutError:
|
|
sys.stderr.write(f"[Phase 2] Resolved handle @{resolved_handle} timed out (30s)\n")
|
|
except Exception as e:
|
|
sys.stderr.write(f"[Phase 2] Resolved handle error: {e}\n")
|
|
|
|
if supplemental_reddit or supplemental_x:
|
|
sys.stderr.write(
|
|
f"[Phase 2] +{len(supplemental_reddit)} Reddit, +{len(supplemental_x)} X\n"
|
|
)
|
|
sys.stderr.flush()
|
|
|
|
return supplemental_reddit, supplemental_x
|
|
|
|
|
|
def run_research(
|
|
topic: str,
|
|
sources: str,
|
|
config: dict,
|
|
selected_models: dict,
|
|
from_date: str,
|
|
to_date: str,
|
|
depth: str = "default",
|
|
mock: bool = False,
|
|
progress: ui.ProgressDisplay = None,
|
|
x_source: str = "xai",
|
|
run_youtube: bool = False,
|
|
run_tiktok: bool = False,
|
|
run_instagram: bool = False,
|
|
run_xiaohongshu: bool = False,
|
|
timeouts: dict = None,
|
|
resolved_handle: str = None,
|
|
do_hackernews: bool = True,
|
|
do_bluesky: bool = True,
|
|
do_truthsocial: bool = True,
|
|
do_polymarket: bool = True,
|
|
no_native_web: bool = False,
|
|
) -> tuple:
|
|
"""Run the research pipeline.
|
|
|
|
Returns:
|
|
Tuple of (reddit_items, x_items, youtube_items, tiktok_items, instagram_items,
|
|
hackernews_items, bluesky_items, truthsocial_items, polymarket_items, web_items, web_needed,
|
|
raw_openai, raw_xai, raw_reddit_enriched,
|
|
reddit_error, x_error, youtube_error, tiktok_error, instagram_error,
|
|
hackernews_error, bluesky_error, truthsocial_error, polymarket_error, web_error)
|
|
|
|
Note: web_needed is True when web search should be performed by the assistant
|
|
(i.e., no native web search API keys are configured). When native web search
|
|
runs, web_items will be populated and web_needed will be False.
|
|
"""
|
|
if timeouts is None:
|
|
timeouts = TIMEOUT_PROFILES[depth]
|
|
future_timeout = timeouts["future"]
|
|
|
|
reddit_items = []
|
|
x_items = []
|
|
youtube_items = []
|
|
tiktok_items = []
|
|
instagram_items = []
|
|
hackernews_items = []
|
|
bluesky_items = []
|
|
truthsocial_items = []
|
|
polymarket_items = []
|
|
web_items = []
|
|
raw_openai = None
|
|
raw_xai = None
|
|
raw_reddit_enriched = []
|
|
reddit_error = None
|
|
x_error = None
|
|
youtube_error = None
|
|
tiktok_error = None
|
|
instagram_error = None
|
|
hackernews_error = None
|
|
bluesky_error = None
|
|
truthsocial_error = None
|
|
polymarket_error = None
|
|
web_error = None
|
|
xiaohongshu_error = None
|
|
|
|
# Determine web search mode
|
|
do_web = sources in ("all", "web", "reddit-web", "x-web")
|
|
web_backend = env.get_web_search_source(config) if (do_web and not no_native_web) else None
|
|
web_needed = do_web and not web_backend
|
|
|
|
# Web-only mode
|
|
if sources == "web":
|
|
if web_backend:
|
|
# Native web search available — run it
|
|
sys.stderr.write(f"[web] Searching via {web_backend}\n")
|
|
sys.stderr.flush()
|
|
try:
|
|
web_items, web_error = _search_web(topic, config, from_date, to_date, depth)
|
|
if web_error and progress:
|
|
progress.show_error(f"Web error: {web_error}")
|
|
except Exception as e:
|
|
web_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"Web error: {e}")
|
|
sys.stderr.write(f"[web] {len(web_items)} results\n")
|
|
sys.stderr.flush()
|
|
else:
|
|
# No native backend — assistant handles WebSearch
|
|
if progress:
|
|
progress.start_web_only()
|
|
progress.end_web_only()
|
|
# Optional Xiaohongshu search in web-only mode.
|
|
if run_xiaohongshu:
|
|
try:
|
|
xhs_items, xiaohongshu_error = _search_xiaohongshu(
|
|
topic, config, from_date, to_date, depth,
|
|
)
|
|
web_items.extend(xhs_items)
|
|
if xiaohongshu_error and progress:
|
|
progress.show_error(f"Xiaohongshu error: {xiaohongshu_error}")
|
|
except Exception as e:
|
|
xiaohongshu_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"Xiaohongshu error: {e}")
|
|
# Still run YouTube/TikTok/Instagram in web-only mode if available
|
|
if run_youtube:
|
|
if progress:
|
|
progress.start_youtube()
|
|
try:
|
|
youtube_items, youtube_error = _search_youtube(topic, from_date, to_date, depth)
|
|
if youtube_error and progress:
|
|
progress.show_error(f"YouTube error: {youtube_error}")
|
|
except Exception as e:
|
|
youtube_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"YouTube error: {e}")
|
|
if progress:
|
|
progress.end_youtube(len(youtube_items))
|
|
if run_tiktok:
|
|
if progress:
|
|
progress.start_tiktok()
|
|
try:
|
|
tiktok_items, tiktok_error = _search_tiktok(topic, from_date, to_date, depth, env.get_tiktok_token(config))
|
|
if tiktok_error and progress:
|
|
progress.show_error(f"TikTok error: {tiktok_error}")
|
|
except Exception as e:
|
|
tiktok_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"TikTok error: {e}")
|
|
if progress:
|
|
progress.end_tiktok(len(tiktok_items))
|
|
if run_instagram:
|
|
if progress:
|
|
progress.start_instagram()
|
|
try:
|
|
ig_timeout = timeouts.get("instagram_future", future_timeout)
|
|
instagram_items, instagram_error = _search_instagram(topic, from_date, to_date, depth, env.get_instagram_token(config))
|
|
if instagram_error and progress:
|
|
progress.show_error(f"Instagram error: {instagram_error}")
|
|
except Exception as e:
|
|
instagram_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"Instagram error: {e}")
|
|
if progress:
|
|
progress.end_instagram(len(instagram_items))
|
|
return reddit_items, x_items, youtube_items, tiktok_items, instagram_items, hackernews_items, bluesky_items, truthsocial_items, polymarket_items, web_items, web_needed, raw_openai, raw_xai, raw_reddit_enriched, reddit_error, x_error, youtube_error, tiktok_error, instagram_error, hackernews_error, bluesky_error, truthsocial_error, polymarket_error, web_error
|
|
|
|
# Determine which searches to run
|
|
do_reddit = sources in ("both", "reddit", "all", "reddit-web")
|
|
do_x = sources in ("both", "x", "all", "x-web")
|
|
# do_hackernews / do_polymarket are always True by default, but can be
|
|
# restricted via the --search flag to run a focused source subset.
|
|
|
|
# Run Reddit, X, YouTube, HN, Polymarket, and Web searches in parallel
|
|
reddit_future = None
|
|
x_future = None
|
|
youtube_future = None
|
|
tiktok_future = None
|
|
instagram_future = None
|
|
xiaohongshu_future = None
|
|
hackernews_future = None
|
|
bluesky_future = None
|
|
truthsocial_future = None
|
|
polymarket_future = None
|
|
web_future = None
|
|
max_workers = (
|
|
2
|
|
+ (1 if run_youtube else 0)
|
|
+ (1 if run_tiktok else 0)
|
|
+ (1 if run_instagram else 0)
|
|
+ (1 if run_xiaohongshu else 0)
|
|
+ (1 if do_hackernews else 0)
|
|
+ (1 if do_bluesky else 0)
|
|
+ (1 if do_truthsocial else 0)
|
|
+ (1 if do_polymarket else 0)
|
|
+ (1 if web_backend else 0)
|
|
)
|
|
|
|
with ThreadPoolExecutor(max_workers=max_workers) as executor:
|
|
# Submit searches
|
|
if do_reddit:
|
|
if progress:
|
|
progress.start_reddit()
|
|
reddit_future = executor.submit(
|
|
_search_reddit, topic, config, selected_models,
|
|
from_date, to_date, depth, mock
|
|
)
|
|
|
|
if do_x:
|
|
if progress:
|
|
progress.start_x()
|
|
x_future = executor.submit(
|
|
_search_x, topic, config, selected_models,
|
|
from_date, to_date, depth, mock, x_source
|
|
)
|
|
|
|
if run_youtube:
|
|
if progress:
|
|
progress.start_youtube()
|
|
youtube_future = executor.submit(
|
|
_search_youtube, topic, from_date, to_date, depth
|
|
)
|
|
|
|
if run_tiktok:
|
|
if progress:
|
|
progress.start_tiktok()
|
|
tiktok_future = executor.submit(
|
|
_search_tiktok, topic, from_date, to_date, depth,
|
|
env.get_tiktok_token(config),
|
|
)
|
|
|
|
if run_instagram:
|
|
if progress:
|
|
progress.start_instagram()
|
|
instagram_future = executor.submit(
|
|
_search_instagram, topic, from_date, to_date, depth,
|
|
env.get_instagram_token(config),
|
|
)
|
|
|
|
if run_xiaohongshu:
|
|
xiaohongshu_future = executor.submit(
|
|
_search_xiaohongshu, topic, config, from_date, to_date, depth,
|
|
)
|
|
|
|
if do_hackernews:
|
|
if progress:
|
|
progress.start_hackernews()
|
|
hackernews_future = executor.submit(
|
|
_search_hackernews, topic, from_date, to_date, depth
|
|
)
|
|
|
|
if do_bluesky:
|
|
bluesky_future = executor.submit(
|
|
_search_bluesky, topic, from_date, to_date, depth, config
|
|
)
|
|
|
|
if do_truthsocial:
|
|
truthsocial_future = executor.submit(
|
|
_search_truthsocial, topic, from_date, to_date, depth, config
|
|
)
|
|
|
|
if do_polymarket:
|
|
if progress:
|
|
progress.start_polymarket()
|
|
polymarket_future = executor.submit(
|
|
_search_polymarket, topic, from_date, to_date, depth
|
|
)
|
|
|
|
if web_backend:
|
|
sys.stderr.write(f"[web] Searching via {web_backend}\n")
|
|
sys.stderr.flush()
|
|
web_future = executor.submit(
|
|
_search_web, topic, config, from_date, to_date, depth
|
|
)
|
|
|
|
# Collect results (with timeouts to prevent indefinite blocking)
|
|
reddit_used_sc = False # Track if ScrapeCreators was used for Reddit
|
|
if reddit_future:
|
|
reddit_timeout = timeouts.get("reddit_future", future_timeout)
|
|
try:
|
|
reddit_items, raw_openai, reddit_error, reddit_used_sc = reddit_future.result(timeout=reddit_timeout)
|
|
if reddit_error and progress:
|
|
progress.show_error(f"Reddit error: {reddit_error}")
|
|
except TimeoutError:
|
|
reddit_error = f"Reddit search timed out after {reddit_timeout}s"
|
|
if progress:
|
|
progress.show_error(reddit_error)
|
|
except Exception as e:
|
|
reddit_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"Reddit error: {e}")
|
|
if progress:
|
|
progress.end_reddit(len(reddit_items))
|
|
|
|
if x_future:
|
|
try:
|
|
x_items, raw_xai, x_error = x_future.result(timeout=future_timeout)
|
|
if x_error and progress:
|
|
progress.show_error(f"X error: {x_error}")
|
|
except TimeoutError:
|
|
x_error = f"X search timed out after {future_timeout}s"
|
|
if progress:
|
|
progress.show_error(x_error)
|
|
except Exception as e:
|
|
x_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"X error: {e}")
|
|
if progress:
|
|
progress.end_x(len(x_items))
|
|
|
|
if youtube_future:
|
|
yt_timeout = timeouts.get("youtube_future", future_timeout)
|
|
try:
|
|
youtube_items, youtube_error = youtube_future.result(timeout=yt_timeout)
|
|
if youtube_error and progress:
|
|
progress.show_error(f"YouTube error: {youtube_error}")
|
|
except TimeoutError:
|
|
youtube_error = f"YouTube search timed out after {yt_timeout}s"
|
|
if progress:
|
|
progress.show_error(youtube_error)
|
|
except Exception as e:
|
|
youtube_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"YouTube error: {e}")
|
|
if progress:
|
|
progress.end_youtube(len(youtube_items))
|
|
|
|
if tiktok_future:
|
|
tk_timeout = timeouts.get("tiktok_future", future_timeout)
|
|
try:
|
|
tiktok_items, tiktok_error = tiktok_future.result(timeout=tk_timeout)
|
|
if tiktok_error and progress:
|
|
progress.show_error(f"TikTok error: {tiktok_error}")
|
|
except TimeoutError:
|
|
tiktok_error = f"TikTok search timed out after {tk_timeout}s"
|
|
if progress:
|
|
progress.show_error(tiktok_error)
|
|
except Exception as e:
|
|
tiktok_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"TikTok error: {e}")
|
|
if progress:
|
|
progress.end_tiktok(len(tiktok_items))
|
|
|
|
if instagram_future:
|
|
ig_timeout = timeouts.get("instagram_future", future_timeout)
|
|
try:
|
|
instagram_items, instagram_error = instagram_future.result(timeout=ig_timeout)
|
|
if instagram_error and progress:
|
|
progress.show_error(f"Instagram error: {instagram_error}")
|
|
except TimeoutError:
|
|
instagram_error = f"Instagram search timed out after {ig_timeout}s"
|
|
if progress:
|
|
progress.show_error(instagram_error)
|
|
except Exception as e:
|
|
instagram_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"Instagram error: {e}")
|
|
if progress:
|
|
progress.end_instagram(len(instagram_items))
|
|
|
|
if xiaohongshu_future:
|
|
try:
|
|
xhs_items, xiaohongshu_error = xiaohongshu_future.result(timeout=future_timeout)
|
|
web_items.extend(xhs_items)
|
|
if xiaohongshu_error and progress:
|
|
progress.show_error(f"Xiaohongshu error: {xiaohongshu_error}")
|
|
except TimeoutError:
|
|
xiaohongshu_error = f"Xiaohongshu search timed out after {future_timeout}s"
|
|
if progress:
|
|
progress.show_error(xiaohongshu_error)
|
|
except Exception as e:
|
|
xiaohongshu_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"Xiaohongshu error: {e}")
|
|
|
|
if hackernews_future:
|
|
hn_timeout = timeouts.get("hackernews_future", future_timeout)
|
|
try:
|
|
hackernews_items, hackernews_error = hackernews_future.result(timeout=hn_timeout)
|
|
if hackernews_error and progress:
|
|
progress.show_error(f"HN error: {hackernews_error}")
|
|
except TimeoutError:
|
|
hackernews_error = f"HN search timed out after {hn_timeout}s"
|
|
if progress:
|
|
progress.show_error(hackernews_error)
|
|
except Exception as e:
|
|
hackernews_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"HN error: {e}")
|
|
if progress:
|
|
progress.end_hackernews(len(hackernews_items))
|
|
|
|
if bluesky_future:
|
|
bsky_timeout = timeouts.get("bluesky_future", future_timeout)
|
|
try:
|
|
bluesky_items, bluesky_error = bluesky_future.result(timeout=bsky_timeout)
|
|
if bluesky_error and progress:
|
|
progress.show_error(f"Bluesky error: {bluesky_error}")
|
|
except TimeoutError:
|
|
bluesky_error = f"Bluesky search timed out after {bsky_timeout}s"
|
|
if progress:
|
|
progress.show_error(bluesky_error)
|
|
except Exception as e:
|
|
bluesky_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"Bluesky error: {e}")
|
|
|
|
if truthsocial_future:
|
|
ts_timeout = timeouts.get("truthsocial_future", future_timeout)
|
|
try:
|
|
truthsocial_items, truthsocial_error = truthsocial_future.result(timeout=ts_timeout)
|
|
if truthsocial_error and progress:
|
|
progress.show_error(f"Truth Social error: {truthsocial_error}")
|
|
except TimeoutError:
|
|
truthsocial_error = f"Truth Social search timed out after {ts_timeout}s"
|
|
if progress:
|
|
progress.show_error(truthsocial_error)
|
|
except Exception as e:
|
|
truthsocial_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"Truth Social error: {e}")
|
|
|
|
if polymarket_future:
|
|
pm_timeout = timeouts.get("polymarket_future", future_timeout)
|
|
try:
|
|
polymarket_items, polymarket_error = polymarket_future.result(timeout=pm_timeout)
|
|
if polymarket_error and progress:
|
|
progress.show_error(f"Polymarket error: {polymarket_error}")
|
|
except TimeoutError:
|
|
polymarket_error = f"Polymarket search timed out after {pm_timeout}s"
|
|
if progress:
|
|
progress.show_error(polymarket_error)
|
|
except Exception as e:
|
|
polymarket_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"Polymarket error: {e}")
|
|
if progress:
|
|
progress.end_polymarket(len(polymarket_items))
|
|
|
|
if web_future:
|
|
try:
|
|
web_items, web_error = web_future.result(timeout=future_timeout)
|
|
if web_error and progress:
|
|
progress.show_error(f"Web error: {web_error}")
|
|
except TimeoutError:
|
|
web_error = f"Web search timed out after {future_timeout}s"
|
|
if progress:
|
|
progress.show_error(web_error)
|
|
except Exception as e:
|
|
web_error = f"{type(e).__name__}: {e}"
|
|
if progress:
|
|
progress.show_error(f"Web error: {e}")
|
|
sys.stderr.write(f"[web] {len(web_items)} results\n")
|
|
sys.stderr.flush()
|
|
|
|
# Enrich Reddit items with real data (parallel, capped)
|
|
# Skip enrichment if ScrapeCreators already provided comments + engagement
|
|
enrich_max = timeouts["enrich_max_items"]
|
|
enrich_total_timeout = timeouts["enrich_total"]
|
|
items_to_enrich = reddit_items[:enrich_max]
|
|
rate_limited = False # Set True if Reddit returns 429 during enrichment
|
|
|
|
if reddit_used_sc and items_to_enrich:
|
|
# ScrapeCreators already enriched items with comments — just copy to raw list
|
|
sys.stderr.write(f"[Reddit] Skipping old enrichment — ScrapeCreators already provided comments\n")
|
|
sys.stderr.flush()
|
|
raw_reddit_enriched = list(reddit_items[:enrich_max])
|
|
items_to_enrich = [] # Skip the enrichment block below
|
|
|
|
if items_to_enrich:
|
|
if progress:
|
|
progress.start_reddit_enrich(1, len(items_to_enrich))
|
|
|
|
if mock:
|
|
# Sequential mock enrichment (fast, no need for parallelism)
|
|
for i, item in enumerate(items_to_enrich):
|
|
if progress and i > 0:
|
|
progress.update_reddit_enrich(i + 1, len(items_to_enrich))
|
|
try:
|
|
mock_thread = load_fixture("reddit_thread_sample.json")
|
|
reddit_items[i] = reddit_enrich.enrich_reddit_item(item, mock_thread)
|
|
except Exception as e:
|
|
if progress:
|
|
progress.show_error(f"Enrich failed for {item.get('url', 'unknown')}: {e}")
|
|
raw_reddit_enriched.append(reddit_items[i])
|
|
else:
|
|
# Parallel enrichment with bounded concurrency and total timeout
|
|
# Uses short HTTP timeout (10s) and 1 retry to fail fast on 429
|
|
completed_count = 0
|
|
rate_limited = False
|
|
with ThreadPoolExecutor(max_workers=5) as enrich_pool:
|
|
futures = {
|
|
enrich_pool.submit(reddit_enrich.enrich_reddit_item, item): i
|
|
for i, item in enumerate(items_to_enrich)
|
|
}
|
|
try:
|
|
for future in as_completed(futures, timeout=enrich_total_timeout):
|
|
idx = futures[future]
|
|
completed_count += 1
|
|
if progress:
|
|
progress.update_reddit_enrich(completed_count, len(items_to_enrich))
|
|
try:
|
|
reddit_items[idx] = future.result(timeout=timeouts["enrich_per"])
|
|
except reddit_enrich.RedditRateLimitError:
|
|
rate_limited = True
|
|
if progress:
|
|
progress.show_error(
|
|
"Reddit rate-limited (429) — skipping remaining enrichment"
|
|
)
|
|
# Cancel remaining futures and bail
|
|
for f in futures:
|
|
f.cancel()
|
|
break
|
|
except Exception as e:
|
|
if progress:
|
|
progress.show_error(
|
|
f"Enrich failed for {items_to_enrich[idx].get('url', 'unknown')}: {e}"
|
|
)
|
|
raw_reddit_enriched.append(reddit_items[idx])
|
|
except TimeoutError:
|
|
if progress:
|
|
progress.show_error(
|
|
f"Enrichment timed out after {enrich_total_timeout}s "
|
|
f"({completed_count}/{len(items_to_enrich)} done)"
|
|
)
|
|
# Keep unenriched items as-is
|
|
for idx in futures.values():
|
|
if reddit_items[idx] not in raw_reddit_enriched:
|
|
raw_reddit_enriched.append(reddit_items[idx])
|
|
|
|
if progress:
|
|
progress.end_reddit_enrich()
|
|
|
|
# Enrich HN stories with comments
|
|
if hackernews_items:
|
|
try:
|
|
hackernews_items = hackernews.enrich_top_stories(hackernews_items, depth=depth)
|
|
except Exception as e:
|
|
sys.stderr.write(f"[HN] Enrichment error: {e}\n")
|
|
sys.stderr.flush()
|
|
|
|
# Phase 2: Supplemental search based on entities from Phase 1
|
|
# Skip on --quick (speed matters), mock mode, or if Reddit is rate-limiting
|
|
# Also skip Reddit supplemental when ScrapeCreators was used (subreddit drilling already done)
|
|
if depth != "quick" and not mock and (reddit_items or x_items):
|
|
sup_reddit, sup_x = _run_supplemental(
|
|
topic, reddit_items, x_items,
|
|
from_date, to_date, depth, x_source, progress,
|
|
skip_reddit=(rate_limited or reddit_used_sc),
|
|
resolved_handle=resolved_handle,
|
|
)
|
|
if sup_reddit:
|
|
reddit_items.extend(sup_reddit)
|
|
if sup_x:
|
|
x_items.extend(sup_x)
|
|
|
|
return reddit_items, x_items, youtube_items, tiktok_items, instagram_items, hackernews_items, bluesky_items, truthsocial_items, polymarket_items, web_items, web_needed, raw_openai, raw_xai, raw_reddit_enriched, reddit_error, x_error, youtube_error, tiktok_error, instagram_error, hackernews_error, bluesky_error, truthsocial_error, polymarket_error, web_error
|
|
|
|
|
|
def main():
|
|
# Fix Unicode output on Windows (cp1252 can't encode emoji)
|
|
if sys.platform == "win32":
|
|
sys.stdout.reconfigure(encoding="utf-8", errors="replace")
|
|
sys.stderr.reconfigure(encoding="utf-8", errors="replace")
|
|
|
|
parser = argparse.ArgumentParser(
|
|
description="Research a topic from the last N days on Reddit + X"
|
|
)
|
|
parser.add_argument("topic", nargs="*", help="Topic to research")
|
|
parser.add_argument("--mock", action="store_true", help="Use fixtures")
|
|
parser.add_argument(
|
|
"--emit",
|
|
choices=["compact", "json", "md", "context", "path"],
|
|
default="compact",
|
|
help="Output mode",
|
|
)
|
|
parser.add_argument(
|
|
"--sources",
|
|
choices=["auto", "reddit", "x", "both"],
|
|
default="auto",
|
|
help="Source selection",
|
|
)
|
|
parser.add_argument(
|
|
"--quick",
|
|
action="store_true",
|
|
help="Faster research with fewer sources (8-12 each)",
|
|
)
|
|
parser.add_argument(
|
|
"--deep",
|
|
action="store_true",
|
|
help="Comprehensive research with more sources (50-70 Reddit, 40-60 X)",
|
|
)
|
|
parser.add_argument(
|
|
"--debug",
|
|
action="store_true",
|
|
help="Enable verbose debug logging",
|
|
)
|
|
parser.add_argument(
|
|
"--include-web",
|
|
action="store_true",
|
|
help="Include general web search alongside Reddit/X (lower weighted)",
|
|
)
|
|
parser.add_argument(
|
|
"--days",
|
|
type=int,
|
|
default=30,
|
|
choices=range(1, 31),
|
|
metavar="N",
|
|
help="Number of days to look back (1-30, default: 30)",
|
|
)
|
|
parser.add_argument(
|
|
"--store",
|
|
action="store_true",
|
|
help="Persist findings to SQLite database (~/.local/share/last30days/research.db)",
|
|
)
|
|
parser.add_argument(
|
|
"--diagnose",
|
|
action="store_true",
|
|
help="Show source availability diagnostics and exit",
|
|
)
|
|
parser.add_argument(
|
|
"--timeout",
|
|
type=int,
|
|
default=None,
|
|
metavar="SECS",
|
|
help="Global timeout in seconds (default: 180, quick: 90, deep: 300)",
|
|
)
|
|
parser.add_argument(
|
|
"--x-handle",
|
|
type=str,
|
|
default=None,
|
|
metavar="HANDLE",
|
|
help="Resolved X handle for topic entity (without @). Searched unfiltered in Phase 2.",
|
|
)
|
|
parser.add_argument(
|
|
"--search",
|
|
type=str,
|
|
default=None,
|
|
metavar="SOURCES",
|
|
help=(
|
|
"Comma-separated list of sources to run. "
|
|
f"Valid: {', '.join(sorted(VALID_SEARCH_SOURCES))}. "
|
|
"Example: --search reddit,hn (default: all configured sources)"
|
|
),
|
|
)
|
|
parser.add_argument(
|
|
"--no-native-web",
|
|
action="store_true",
|
|
default=False,
|
|
help="Skip native web search backends (Parallel/Brave/OpenRouter). Use when the assistant has its own WebSearch tool.",
|
|
)
|
|
parser.add_argument(
|
|
"--save-dir",
|
|
type=str,
|
|
default=None,
|
|
metavar="DIR",
|
|
help="Auto-save raw research output to DIR/{topic-slug}.md",
|
|
)
|
|
|
|
args = parser.parse_args()
|
|
args.topic = " ".join(args.topic) if args.topic else None
|
|
|
|
# Enable debug logging if requested
|
|
if args.debug:
|
|
os.environ["LAST30DAYS_DEBUG"] = "1"
|
|
# Re-import http to pick up debug flag
|
|
from lib import http as http_module
|
|
http_module.DEBUG = True
|
|
|
|
# Determine depth
|
|
if args.quick and args.deep:
|
|
print("Error: Cannot use both --quick and --deep", file=sys.stderr)
|
|
sys.exit(1)
|
|
elif args.quick:
|
|
depth = "quick"
|
|
elif args.deep:
|
|
depth = "deep"
|
|
else:
|
|
depth = "default"
|
|
|
|
# Install global timeout watchdog
|
|
timeouts = TIMEOUT_PROFILES[depth]
|
|
global_timeout = args.timeout or timeouts["global"]
|
|
_install_global_timeout(global_timeout)
|
|
|
|
# Load config
|
|
config = env.get_config()
|
|
|
|
# Detect first run (no SETUP_COMPLETE in config)
|
|
first_run = setup_wizard.is_first_run(config)
|
|
|
|
# On first run, block Bird's Node.js sweet-cookie scanner from probing
|
|
# browser cookies before the user has given consent via the setup wizard.
|
|
# Explicit AUTH_TOKEN (from env var) is fine — only block browser scanning.
|
|
if first_run and config.get('_AUTH_TOKEN_SOURCE') != 'env':
|
|
os.environ['BIRD_DISABLE_BROWSER_COOKIES'] = '1'
|
|
|
|
# Inject .env credentials into Bird module before auth check.
|
|
# On first run (no SETUP_COMPLETE), only inject explicit env var credentials —
|
|
# skip if AUTH_TOKEN came from browser cookies (no consent yet).
|
|
if first_run:
|
|
auth_source = config.get('_AUTH_TOKEN_SOURCE')
|
|
if auth_source == 'env':
|
|
bird_x.set_credentials(config.get('AUTH_TOKEN'), config.get('CT0'))
|
|
else:
|
|
bird_x.set_credentials(config.get('AUTH_TOKEN'), config.get('CT0'))
|
|
|
|
# Auto-detect Bird (no prompts - just use it if available)
|
|
x_source_status = env.get_x_source_status(config)
|
|
x_source = x_source_status["source"] # 'bird', 'xai', or None
|
|
x_method = x_source_status.get("method") # 'env', 'browser-firefox', 'api', etc.
|
|
|
|
# Auto-detect yt-dlp for YouTube search
|
|
has_ytdlp = env.is_ytdlp_available()
|
|
|
|
# Auto-detect ScrapeCreators/Apify for TikTok
|
|
has_tiktok = env.is_tiktok_available(config)
|
|
|
|
# Auto-detect ScrapeCreators for Instagram
|
|
has_instagram = env.is_instagram_available(config)
|
|
|
|
# Auto-detect Xiaohongshu HTTP API (requires service + login)
|
|
has_xiaohongshu = env.is_xiaohongshu_available(config)
|
|
|
|
# Auto-detect Bluesky (requires BSKY_HANDLE + BSKY_APP_PASSWORD)
|
|
has_bluesky = env.is_bluesky_available(config)
|
|
|
|
# Auto-detect Truth Social (requires TRUTHSOCIAL_TOKEN)
|
|
has_truthsocial = env.is_truthsocial_available(config)
|
|
|
|
# --diagnose: show source availability and exit
|
|
if args.diagnose:
|
|
web_source = env.get_web_search_source(config)
|
|
diag = {
|
|
"openai": bool(config.get("OPENAI_API_KEY")),
|
|
"reddit_public": True,
|
|
"xai": bool(config.get("XAI_API_KEY")),
|
|
"x_source": x_source_status["source"],
|
|
"x_method": x_source_status.get("method"),
|
|
"bird_installed": x_source_status["bird_installed"],
|
|
"bird_authenticated": x_source_status["bird_authenticated"],
|
|
"bird_username": x_source_status.get("bird_username"),
|
|
"youtube": has_ytdlp,
|
|
"tiktok": has_tiktok,
|
|
"instagram": has_instagram,
|
|
"xiaohongshu": has_xiaohongshu,
|
|
"xiaohongshu_api_base": env.get_xiaohongshu_api_base(config),
|
|
"hackernews": True,
|
|
"bluesky": has_bluesky,
|
|
"truthsocial": has_truthsocial,
|
|
"polymarket": True,
|
|
"web_search_backend": web_source,
|
|
"exa": bool(config.get("EXA_API_KEY")),
|
|
"parallel_ai": bool(config.get("PARALLEL_API_KEY")),
|
|
"brave": bool(config.get("BRAVE_API_KEY")),
|
|
"openrouter": bool(config.get("OPENROUTER_API_KEY")),
|
|
}
|
|
print(json.dumps(diag, indent=2))
|
|
sys.exit(0)
|
|
|
|
# Handle 'setup' subcommand
|
|
if args.topic and args.topic.strip().lower() == "setup":
|
|
results = setup_wizard.run_auto_setup(config)
|
|
# Write config
|
|
env_path = env.CONFIG_FILE
|
|
if env_path:
|
|
written = setup_wizard.write_setup_config(env_path)
|
|
results["env_written"] = written
|
|
else:
|
|
results["env_written"] = False
|
|
print(setup_wizard.get_setup_status_text(results))
|
|
sys.exit(0)
|
|
|
|
# Validate topic (--diagnose doesn't need one)
|
|
if not args.topic:
|
|
print("Error: Please provide a topic to research.", file=sys.stderr)
|
|
print("Usage: python3 last30days.py <topic> [options]", file=sys.stderr)
|
|
sys.exit(1)
|
|
|
|
# Initialize progress display with topic
|
|
progress = ui.ProgressDisplay(args.topic, show_banner=True)
|
|
|
|
# Show status banner (free-first design — lead with what works)
|
|
web_source = env.get_web_search_source(config)
|
|
reddit_source = env.get_reddit_source(config)
|
|
diag = {
|
|
"setup_complete": bool(config.get("SETUP_COMPLETE")),
|
|
"reddit_source": reddit_source, # 'scrapecreators', 'openai', or None
|
|
"x_source": x_source_status["source"],
|
|
"x_method": x_source_status.get("method"),
|
|
"youtube": has_ytdlp,
|
|
"tiktok": has_tiktok,
|
|
"instagram": has_instagram,
|
|
"hackernews": True,
|
|
"polymarket": True,
|
|
"bluesky": has_bluesky,
|
|
"truthsocial": has_truthsocial,
|
|
"xiaohongshu": has_xiaohongshu,
|
|
"scrapecreators": bool(config.get("SCRAPECREATORS_API_KEY")),
|
|
"web_search_backend": "deferred to assistant" if args.no_native_web else web_source,
|
|
}
|
|
ui.show_diagnostic_banner(diag)
|
|
|
|
# Check available sources (now accounts for Bird/cookie auth automatically)
|
|
available = env.get_available_sources(config)
|
|
|
|
# Mock mode can work without keys
|
|
if args.mock:
|
|
if args.sources == "auto":
|
|
sources = "both"
|
|
else:
|
|
sources = args.sources
|
|
else:
|
|
# Validate requested sources against available
|
|
sources, error = env.validate_sources(args.sources, available, args.include_web)
|
|
if error:
|
|
# If it's a warning about WebSearch fallback, print but continue
|
|
if "WebSearch fallback" in error:
|
|
print(f"Note: {error}", file=sys.stderr)
|
|
else:
|
|
print(f"Error: {error}", file=sys.stderr)
|
|
sys.exit(1)
|
|
|
|
# Get date range
|
|
from_date, to_date = dates.get_date_range(args.days)
|
|
|
|
# Check what keys are missing for promo messaging
|
|
missing_keys = env.get_missing_keys(config)
|
|
|
|
# Show NUX / promo for missing keys BEFORE research
|
|
if missing_keys != 'none':
|
|
progress.show_promo(missing_keys, diag=diag)
|
|
|
|
# Select models
|
|
if args.mock:
|
|
# Use mock models
|
|
mock_openai_models = load_fixture("models_openai_sample.json").get("data", [])
|
|
mock_xai_models = load_fixture("models_xai_sample.json").get("data", [])
|
|
selected_models = models.get_models(
|
|
{
|
|
"OPENAI_API_KEY": "mock",
|
|
"XAI_API_KEY": "mock",
|
|
**config,
|
|
},
|
|
mock_openai_models,
|
|
mock_xai_models,
|
|
)
|
|
else:
|
|
selected_models = models.get_models(config)
|
|
|
|
# Determine mode string
|
|
if sources == "all":
|
|
mode = "all" # reddit + x + web
|
|
elif sources == "both":
|
|
mode = "both" # reddit + x
|
|
elif sources == "reddit":
|
|
mode = "reddit-only"
|
|
elif sources == "reddit-web":
|
|
mode = "reddit-web"
|
|
elif sources == "x":
|
|
mode = "x-only"
|
|
elif sources == "x-web":
|
|
mode = "x-web"
|
|
elif sources == "web":
|
|
mode = "web-only"
|
|
else:
|
|
mode = sources
|
|
|
|
# Detect query type for source tiering and scoring adjustments
|
|
query_type = qt.detect_query_type(args.topic)
|
|
|
|
# Apply --search flag: restrict sources to the specified subset
|
|
# Source defaults are query-type-aware (Truth Social always opt-in,
|
|
# Bluesky only for query types where it adds signal)
|
|
search_do_hackernews = qt.is_source_enabled("hn", query_type) if not args.search else True
|
|
search_do_bluesky = has_bluesky and qt.is_source_enabled("bluesky", query_type)
|
|
search_do_truthsocial = False # Always opt-in (requires --search truthsocial)
|
|
search_do_polymarket = qt.is_source_enabled("polymarket", query_type)
|
|
search_run_youtube = has_ytdlp and qt.is_source_enabled("youtube", query_type)
|
|
search_run_tiktok = has_tiktok and qt.is_source_enabled("tiktok", query_type)
|
|
search_run_instagram = has_instagram and qt.is_source_enabled("instagram", query_type)
|
|
search_run_xiaohongshu = has_xiaohongshu
|
|
if args.search:
|
|
search_sources = parse_search_flag(args.search)
|
|
has_reddit = "reddit" in search_sources
|
|
has_x = "x" in search_sources
|
|
search_do_hackernews = "hn" in search_sources
|
|
search_do_bluesky = ("bluesky" in search_sources or "bsky" in search_sources) and has_bluesky
|
|
search_do_truthsocial = ("truthsocial" in search_sources or "truth" in search_sources) and has_truthsocial
|
|
search_do_polymarket = "polymarket" in search_sources
|
|
search_run_youtube = "youtube" in search_sources and has_ytdlp
|
|
search_run_tiktok = "tiktok" in search_sources and has_tiktok
|
|
search_run_instagram = "instagram" in search_sources and has_instagram
|
|
# If explicitly requested, attempt Xiaohongshu even when preflight says unavailable.
|
|
search_run_xiaohongshu = "xiaohongshu" in search_sources
|
|
include_search_web = "web" in search_sources
|
|
# Map to existing sources string
|
|
if has_reddit and has_x:
|
|
sources = "both" + ("-web" if include_search_web else "")
|
|
sources = "all" if include_search_web else "both"
|
|
elif has_reddit:
|
|
sources = "reddit-web" if include_search_web else "reddit"
|
|
elif has_x:
|
|
sources = "x-web" if include_search_web else "x"
|
|
else:
|
|
sources = "web" # hn/polymarket only; no Reddit/X
|
|
|
|
# Run research
|
|
reddit_items, x_items, youtube_items, tiktok_items, instagram_items, hackernews_items, bluesky_items, truthsocial_items, polymarket_items, web_items, web_needed, raw_openai, raw_xai, raw_reddit_enriched, reddit_error, x_error, youtube_error, tiktok_error, instagram_error, hackernews_error, bluesky_error, truthsocial_error, polymarket_error, web_error = run_research(
|
|
args.topic,
|
|
sources,
|
|
config,
|
|
selected_models,
|
|
from_date,
|
|
to_date,
|
|
depth,
|
|
args.mock,
|
|
progress,
|
|
x_source=x_source or "xai",
|
|
run_youtube=search_run_youtube,
|
|
run_tiktok=search_run_tiktok,
|
|
run_instagram=search_run_instagram,
|
|
run_xiaohongshu=search_run_xiaohongshu,
|
|
timeouts=timeouts,
|
|
resolved_handle=args.x_handle,
|
|
do_hackernews=search_do_hackernews,
|
|
do_bluesky=search_do_bluesky,
|
|
do_truthsocial=search_do_truthsocial,
|
|
do_polymarket=search_do_polymarket,
|
|
no_native_web=args.no_native_web,
|
|
)
|
|
|
|
# Processing phase
|
|
progress.start_processing()
|
|
|
|
# Normalize items
|
|
normalized_reddit = normalize.normalize_reddit_items(reddit_items, from_date, to_date)
|
|
normalized_x = normalize.normalize_x_items(x_items, from_date, to_date)
|
|
normalized_youtube = normalize.normalize_youtube_items(youtube_items, from_date, to_date) if youtube_items else []
|
|
normalized_tiktok = normalize.normalize_tiktok_items(tiktok_items, from_date, to_date) if tiktok_items else []
|
|
normalized_ig = normalize.normalize_instagram_items(instagram_items, from_date, to_date) if instagram_items else []
|
|
normalized_hn = normalize.normalize_hackernews_items(hackernews_items, from_date, to_date) if hackernews_items else []
|
|
normalized_bsky = normalize.normalize_bluesky_items(bluesky_items, from_date, to_date) if bluesky_items else []
|
|
normalized_ts = normalize.normalize_truthsocial_items(truthsocial_items, from_date, to_date) if truthsocial_items else []
|
|
normalized_pm = normalize.normalize_polymarket_items(polymarket_items, from_date, to_date) if polymarket_items else []
|
|
normalized_web = websearch.normalize_websearch_items(web_items, from_date, to_date) if web_items else []
|
|
|
|
# Hard date filter: exclude items with verified dates outside the range
|
|
# This is the safety net - even if prompts let old content through, this filters it
|
|
filtered_reddit = normalize.filter_by_date_range(normalized_reddit, from_date, to_date)
|
|
filtered_x = normalize.filter_by_date_range(normalized_x, from_date, to_date)
|
|
# YouTube: skip hard date filter — youtube_yt.py already applies a soft filter
|
|
# that prefers recent videos but keeps older ones for evergreen topics.
|
|
# YouTube content has a longer shelf life than tweets/posts.
|
|
filtered_youtube = normalized_youtube
|
|
# TikTok: hard date filter (tiktok.py already pre-filters, but safety net)
|
|
filtered_tiktok = normalize.filter_by_date_range(normalized_tiktok, from_date, to_date) if normalized_tiktok else []
|
|
# Instagram: hard date filter (instagram.py already pre-filters, but safety net)
|
|
filtered_ig = normalize.filter_by_date_range(normalized_ig, from_date, to_date) if normalized_ig else []
|
|
filtered_hn = normalize.filter_by_date_range(normalized_hn, from_date, to_date) if normalized_hn else []
|
|
filtered_bsky = normalize.filter_by_date_range(normalized_bsky, from_date, to_date) if normalized_bsky else []
|
|
filtered_ts = normalize.filter_by_date_range(normalized_ts, from_date, to_date) if normalized_ts else []
|
|
# Polymarket: skip hard date filter - markets are active/traded, updatedAt is fine
|
|
filtered_pm = normalized_pm
|
|
filtered_web = normalize.filter_by_date_range(normalized_web, from_date, to_date) if normalized_web else []
|
|
|
|
# Score items
|
|
scored_reddit = score.score_reddit_items(filtered_reddit)
|
|
scored_x = score.score_x_items(filtered_x)
|
|
scored_youtube = score.score_youtube_items(filtered_youtube) if filtered_youtube else []
|
|
scored_tiktok = score.score_tiktok_items(filtered_tiktok) if filtered_tiktok else []
|
|
scored_ig = score.score_instagram_items(filtered_ig) if filtered_ig else []
|
|
scored_hn = score.score_hackernews_items(filtered_hn) if filtered_hn else []
|
|
scored_bsky = score.score_bluesky_items(filtered_bsky) if filtered_bsky else []
|
|
scored_ts = score.score_truthsocial_items(filtered_ts) if filtered_ts else []
|
|
scored_pm = score.score_polymarket_items(filtered_pm) if filtered_pm else []
|
|
scored_web = score.score_websearch_items(filtered_web, query_type=query_type) if filtered_web else []
|
|
|
|
# Sort items (query-type-aware tiebreaker ordering)
|
|
sorted_reddit = score.sort_items(scored_reddit, query_type=query_type)
|
|
sorted_x = score.sort_items(scored_x, query_type=query_type)
|
|
sorted_youtube = score.sort_items(scored_youtube, query_type=query_type) if scored_youtube else []
|
|
sorted_tiktok = score.sort_items(scored_tiktok, query_type=query_type) if scored_tiktok else []
|
|
sorted_ig = score.sort_items(scored_ig, query_type=query_type) if scored_ig else []
|
|
sorted_hn = score.sort_items(scored_hn, query_type=query_type) if scored_hn else []
|
|
sorted_bsky = score.sort_items(scored_bsky, query_type=query_type) if scored_bsky else []
|
|
sorted_ts = score.sort_items(scored_ts, query_type=query_type) if scored_ts else []
|
|
sorted_pm = score.sort_items(scored_pm, query_type=query_type) if scored_pm else []
|
|
sorted_web = score.sort_items(scored_web, query_type=query_type) if scored_web else []
|
|
|
|
# Dedupe items
|
|
deduped_reddit = dedupe.dedupe_reddit(sorted_reddit)
|
|
deduped_x = dedupe.dedupe_x(sorted_x)
|
|
deduped_youtube = dedupe.dedupe_youtube(sorted_youtube) if sorted_youtube else []
|
|
deduped_tiktok = dedupe.dedupe_tiktok(sorted_tiktok) if sorted_tiktok else []
|
|
deduped_ig = dedupe.dedupe_instagram(sorted_ig) if sorted_ig else []
|
|
deduped_hn = dedupe.dedupe_hackernews(sorted_hn) if sorted_hn else []
|
|
deduped_bsky = dedupe.dedupe_bluesky(sorted_bsky) if sorted_bsky else []
|
|
deduped_ts = dedupe.dedupe_truthsocial(sorted_ts) if sorted_ts else []
|
|
deduped_pm = dedupe.dedupe_polymarket(sorted_pm) if sorted_pm else []
|
|
deduped_web = websearch.dedupe_websearch(sorted_web) if sorted_web else []
|
|
|
|
# Post-retrieval relevance filter: drop low-relevance items per source
|
|
deduped_reddit = score.relevance_filter(deduped_reddit, "REDDIT")
|
|
deduped_x = score.relevance_filter(deduped_x, "X")
|
|
deduped_youtube = score.relevance_filter(deduped_youtube, "YOUTUBE")
|
|
deduped_tiktok = score.relevance_filter(deduped_tiktok, "TIKTOK")
|
|
deduped_ig = score.relevance_filter(deduped_ig, "INSTAGRAM")
|
|
deduped_hn = score.relevance_filter(deduped_hn, "HN")
|
|
deduped_bsky = score.relevance_filter(deduped_bsky, "BLUESKY")
|
|
deduped_ts = score.relevance_filter(deduped_ts, "TRUTHSOCIAL")
|
|
deduped_pm = score.relevance_filter(deduped_pm, "POLYMARKET") if deduped_pm else []
|
|
|
|
# Cross-source linking: annotate items that discuss the same story
|
|
dedupe.cross_source_link(
|
|
deduped_reddit, deduped_x, deduped_youtube, deduped_tiktok, deduped_ig, deduped_hn, deduped_bsky, deduped_ts, deduped_pm, deduped_web,
|
|
)
|
|
|
|
progress.end_processing()
|
|
|
|
# Create report
|
|
report = schema.create_report(
|
|
args.topic,
|
|
from_date,
|
|
to_date,
|
|
mode,
|
|
selected_models.get("openai"),
|
|
selected_models.get("xai"),
|
|
)
|
|
report.reddit = deduped_reddit
|
|
report.x = deduped_x
|
|
report.youtube = deduped_youtube
|
|
report.tiktok = deduped_tiktok
|
|
report.instagram = deduped_ig
|
|
report.hackernews = deduped_hn
|
|
report.bluesky = deduped_bsky
|
|
report.truthsocial = deduped_ts
|
|
report.polymarket = deduped_pm
|
|
report.web = deduped_web
|
|
report.reddit_error = reddit_error
|
|
report.x_error = x_error
|
|
report.youtube_error = youtube_error
|
|
report.tiktok_error = tiktok_error
|
|
report.instagram_error = instagram_error
|
|
report.hackernews_error = hackernews_error
|
|
report.bluesky_error = bluesky_error
|
|
report.truthsocial_error = truthsocial_error
|
|
report.polymarket_error = polymarket_error
|
|
report.web_error = web_error
|
|
report.resolved_x_handle = args.x_handle
|
|
|
|
# Generate context snippet
|
|
report.context_snippet_md = render.render_context_snippet(report)
|
|
|
|
# Write outputs
|
|
render.write_outputs(report, raw_openai, raw_xai, raw_reddit_enriched)
|
|
|
|
# Show completion
|
|
if sources == "web":
|
|
progress.show_web_only_complete()
|
|
else:
|
|
progress.show_complete(len(deduped_reddit), len(deduped_x), len(deduped_youtube), len(deduped_hn), len(deduped_pm), len(deduped_tiktok), len(deduped_ig))
|
|
|
|
# Build source info for status footer
|
|
source_info = {}
|
|
if not x_source:
|
|
if x_source_status["bird_installed"]:
|
|
source_info["x_skip_reason"] = "Bird installed but not authenticated — log into x.com in browser"
|
|
else:
|
|
source_info["x_skip_reason"] = "No Bird CLI, XAI_API_KEY, or SCRAPECREATORS_API_KEY"
|
|
if not has_ytdlp:
|
|
source_info["youtube_skip_reason"] = "yt-dlp not installed — fix: brew install yt-dlp"
|
|
elif has_ytdlp and not report.youtube:
|
|
source_info["youtube_skip_reason"] = "0 results (query may be too specific)"
|
|
if not has_tiktok:
|
|
source_info["tiktok_skip_reason"] = "No SCRAPECREATORS_API_KEY - sign up at scrapecreators.com (100 free API calls, no credit card)"
|
|
if not has_instagram:
|
|
source_info["instagram_skip_reason"] = "No SCRAPECREATORS_API_KEY - sign up at scrapecreators.com (100 free API calls, no credit card)"
|
|
if not has_xiaohongshu:
|
|
source_info["xiaohongshu_skip_reason"] = (
|
|
f"Xiaohongshu API unavailable or not logged in - start xiaohongshu-mcp and login "
|
|
f"(base: {env.get_xiaohongshu_api_base(config)})"
|
|
)
|
|
if not web_source:
|
|
source_info["web_skip_reason"] = "assistant will use WebSearch (add BRAVE_API_KEY for native search)"
|
|
|
|
# Compute quality score and upgrade nudge
|
|
research_results = {
|
|
"x_error": x_error,
|
|
"youtube_error": youtube_error,
|
|
"reddit_error": reddit_error,
|
|
}
|
|
quality = quality_nudge.compute_quality_score(config, research_results)
|
|
|
|
# Output result
|
|
output_result(report, args.emit, web_needed, args.topic, from_date, to_date, missing_keys, args.days, source_info, first_run=first_run, quality=quality)
|
|
|
|
# Auto-save raw research to file if --save-dir is set
|
|
if args.save_dir:
|
|
import re
|
|
save_dir = Path(args.save_dir).expanduser()
|
|
save_dir.mkdir(parents=True, exist_ok=True)
|
|
slug = re.sub(r'[^a-z0-9]+', '-', args.topic.lower()).strip('-')[:60]
|
|
save_path = save_dir / f"{slug}-raw.md"
|
|
if save_path.exists():
|
|
save_path = save_dir / f"{slug}-raw-{datetime.now().strftime('%Y-%m-%d')}.md"
|
|
content = render.render_compact(report, missing_keys=missing_keys)
|
|
if quality and quality.get("nudge_text"):
|
|
content += "\n" + render.render_quality_nudge(quality)
|
|
content += "\n" + render.render_source_status(report, source_info)
|
|
save_path.write_text(content, encoding="utf-8")
|
|
print(f"📎 {save_path}", file=sys.stderr)
|
|
|
|
# Persist findings to SQLite if requested
|
|
if args.store:
|
|
import store as store_mod
|
|
store_mod.init_db()
|
|
topic_row = store_mod.add_topic(args.topic)
|
|
topic_id = topic_row["id"]
|
|
run_id = store_mod.record_run(topic_id, source_mode=mode, status="completed")
|
|
|
|
findings = []
|
|
for item in deduped_reddit:
|
|
findings.append({
|
|
"source": "reddit",
|
|
"url": item.url,
|
|
"title": item.title,
|
|
"author": item.subreddit,
|
|
"content": item.title,
|
|
"engagement_score": item.engagement.score if item.engagement else 0,
|
|
"relevance_score": item.relevance,
|
|
})
|
|
for item in deduped_x:
|
|
findings.append({
|
|
"source": "x",
|
|
"url": item.url,
|
|
"title": item.text[:100],
|
|
"author": item.author_handle,
|
|
"content": item.text,
|
|
"engagement_score": item.engagement.likes if item.engagement else 0,
|
|
"relevance_score": item.relevance,
|
|
})
|
|
for item in deduped_youtube:
|
|
findings.append({
|
|
"source": "youtube",
|
|
"url": item.url,
|
|
"title": item.title,
|
|
"author": item.channel_name,
|
|
"content": item.transcript_snippet[:500] if item.transcript_snippet else item.title,
|
|
"engagement_score": item.engagement.views if item.engagement and item.engagement.views else 0,
|
|
"relevance_score": item.relevance,
|
|
})
|
|
for item in deduped_hn:
|
|
findings.append({
|
|
"source": "hackernews",
|
|
"url": item.hn_url,
|
|
"title": item.title,
|
|
"author": item.author,
|
|
"content": item.title,
|
|
"engagement_score": item.engagement.score if item.engagement else 0,
|
|
"relevance_score": item.relevance,
|
|
})
|
|
for item in deduped_bsky:
|
|
findings.append({
|
|
"source": "bluesky",
|
|
"url": item.url,
|
|
"title": item.text[:100],
|
|
"author": item.author_handle,
|
|
"content": item.text,
|
|
"engagement_score": item.engagement.likes if item.engagement else 0,
|
|
"relevance_score": item.relevance,
|
|
})
|
|
for item in deduped_pm:
|
|
findings.append({
|
|
"source": "polymarket",
|
|
"url": item.url,
|
|
"title": item.question,
|
|
"author": "polymarket",
|
|
"content": item.title,
|
|
"engagement_score": item.engagement.volume if item.engagement and item.engagement.volume else 0,
|
|
"relevance_score": item.relevance,
|
|
})
|
|
for item in deduped_ig:
|
|
findings.append({
|
|
"source": "instagram",
|
|
"url": item.url,
|
|
"title": item.text[:100],
|
|
"author": item.author_name,
|
|
"content": item.caption_snippet[:500] if item.caption_snippet else item.text,
|
|
"engagement_score": item.engagement.views if item.engagement and item.engagement.views else 0,
|
|
"relevance_score": item.relevance,
|
|
})
|
|
for item in deduped_web:
|
|
findings.append({
|
|
"source": "web",
|
|
"url": item.url,
|
|
"title": item.title,
|
|
"author": item.source_domain,
|
|
"content": item.snippet,
|
|
"engagement_score": 0,
|
|
"relevance_score": item.relevance,
|
|
})
|
|
|
|
counts = store_mod.store_findings(run_id, topic_id, findings)
|
|
store_mod.update_run(
|
|
run_id,
|
|
status="completed",
|
|
findings_new=counts["new"],
|
|
findings_updated=counts["updated"],
|
|
)
|
|
sys.stderr.write(
|
|
f"[store] Saved {counts['new']} new, {counts['updated']} updated findings\n"
|
|
)
|
|
sys.stderr.flush()
|
|
|
|
|
|
def output_result(
|
|
report: schema.Report,
|
|
emit_mode: str,
|
|
web_needed: bool = False,
|
|
topic: str = "",
|
|
from_date: str = "",
|
|
to_date: str = "",
|
|
missing_keys: str = "none",
|
|
days: int = 30,
|
|
source_info: dict = None,
|
|
first_run: bool = False,
|
|
quality: dict = None,
|
|
):
|
|
"""Output the result based on emit mode."""
|
|
if emit_mode == "compact":
|
|
# Emit first-run flag before research output so SKILL.md can detect it
|
|
if first_run:
|
|
print("FIRST_RUN: true")
|
|
print("")
|
|
print(render.render_compact(report, missing_keys=missing_keys))
|
|
# Quality nudge (right before source status/stats block)
|
|
if quality and quality.get("nudge_text"):
|
|
print(render.render_quality_nudge(quality))
|
|
# Append source status footer
|
|
print(render.render_source_status(report, source_info))
|
|
elif emit_mode == "json":
|
|
print(json.dumps(report.to_dict(), indent=2))
|
|
elif emit_mode == "md":
|
|
print(render.render_full_report(report))
|
|
elif emit_mode == "context":
|
|
print(report.context_snippet_md)
|
|
elif emit_mode == "path":
|
|
print(render.get_context_path())
|
|
|
|
# Output WebSearch instructions if needed
|
|
if web_needed:
|
|
print("\n" + "="*60)
|
|
print("### WEBSEARCH REQUIRED ###")
|
|
print("="*60)
|
|
print(f"Topic: {topic}")
|
|
print(f"Date range: {from_date} to {to_date}")
|
|
print("")
|
|
print("Assistant: Use your web search tool to find 8-15 relevant web pages.")
|
|
print("EXCLUDE: reddit.com, x.com, twitter.com (already covered above)")
|
|
print(f"INCLUDE: blogs, docs, news, tutorials from the last {days} days")
|
|
print("")
|
|
print("After searching, synthesize WebSearch results WITH the Reddit/X")
|
|
print("results above. WebSearch items should rank LOWER than comparable")
|
|
print("Reddit/X items (they lack engagement metrics).")
|
|
print("="*60)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|