返回 last30days-skill
last30days.py
根目录 / skills / last30days / scripts / last30days.py
1 #!/usr/bin/env python3
2 # fmt: off
3 # ruff: noqa: E402
4 """last30days CLI."""
5
6 from __future__ import annotations
7
8 import argparse
9 import atexit
10 import datetime
11 import hashlib
12 import json
13 import os
14 import re
15 import signal
16 import sqlite3
17 import sys
18 from collections.abc import Callable
19 from pathlib import Path
20
21 MIN_PYTHON = (3, 12)
22
23
24 def ensure_supported_python(version_info: tuple[int, int, int] | object | None = None) -> None:
25 if version_info is None:
26 version_info = sys.version_info
27 major, minor, micro = tuple(version_info[:3])
28 if (major, minor) >= MIN_PYTHON:
29 return
30 req = f"{MIN_PYTHON[0]}.{MIN_PYTHON[1]}"
31 sys.stderr.write(
32 f"last30days v3 requires Python {req}+.\n"
33 f"Detected Python {major}.{minor}.{micro}.\n"
34 f"Install with:\n"
35 f" Mac: brew install python@{req}\n"
36 f" Windows: winget install Python.Python.{req}\n"
37 f" Linux: sudo apt install python{req} (or pyenv install {req})\n"
38 f"Then rerun: python{req} <path-to-script> setup\n"
39 )
40 raise SystemExit(1)
41
42
43 ensure_supported_python()
44
45 if os.name == "nt":
46 for stream in (sys.stdout, sys.stderr):
47 if hasattr(stream, "reconfigure"):
48 stream.reconfigure(encoding="utf-8", errors="replace")
49
50 SCRIPT_DIR = Path(__file__).parent.resolve()
51 sys.path.insert(0, str(SCRIPT_DIR))
52
53 from lib import competitors as competitors_mod, corpus, dates, discovery_handoff, env, freshness, html_render, http, permission_preflight, pipeline, reddit, registers, render, schema, subproc, ui, x_envelope
54
55 atexit.register(subproc.cleanup_children)
56
57
58 def _on_sigterm(signum, frame) -> None:
59 """SIGTERM handler: clean descendant groups, then die as SIGTERM.
60
61 Every run_with_timeout child runs in its own pgid (lib/subproc.py via
62 os.setsid), so a group kill aimed at the engine can never reach them;
63 only the lib.subproc registry can. atexit never runs on a signal death,
64 so without this handler an MCP timeout would orphan node bird-search,
65 yt-dlp, and the digg CLI. The MCP server SIGTERMs the engine group
66 first, giving this handler room to killpg() each registered child
67 group before the SIGKILL backstop. Restoring the default disposition
68 and re-raising preserves killed-by-SIGTERM semantics for the parent.
69 """
70 subproc.cleanup_children()
71 signal.signal(signal.SIGTERM, signal.SIG_DFL)
72 os.kill(os.getpid(), signal.SIGTERM)
73
74
75 def _install_sigterm_handler() -> None:
76 try:
77 signal.signal(signal.SIGTERM, _on_sigterm)
78 except (ValueError, OSError, RuntimeError):
79 pass
80
81
82 def parse_meta_ads_page(raw: str) -> str:
83 """Extract an Ad Library page id from a flag value, or "" if there is none.
84
85 Accepts a bare numeric id or any Ad Library URL carrying
86 ``view_all_page_id``. A ``facebook.com/<vanity>`` URL is deliberately
87 rejected rather than guessed at: a vanity handle is not a page id, and one
88 live check resolved a brand-looking handle to a private person's profile.
89 """
90 value = str(raw or "").strip()
91 if not value:
92 return ""
93 if re.fullmatch(r"\d{5,20}", value):
94 return value
95 match = re.search(r"view_all_page_id=(\d{5,20})", value)
96 return match.group(1) if match else ""
97
98
99 def parse_search_flag(raw: str, flag_name: str = "--search") -> list[str]:
100 sources = []
101 for source in raw.split(","):
102 source = source.strip().lower()
103 if not source:
104 continue
105 normalized = pipeline.SEARCH_ALIAS.get(source, source)
106 if normalized not in pipeline.MOCK_AVAILABLE_SOURCES:
107 raise SystemExit(f"Unknown search source in {flag_name}: {source}")
108 if normalized not in sources:
109 sources.append(normalized)
110 if not sources:
111 raise SystemExit(f"{flag_name} requires at least one source.")
112 return sources
113
114 def parse_as_of_date_arg(value: str) -> str:
115 try:
116 parsed = dates.parse_as_of_date(value)
117 except ValueError as exc:
118 raise argparse.ArgumentTypeError(str(exc)) from exc
119 return parsed
120
121 def resolve_requested_sources(args_search: str | None, config: dict) -> list[str] | None:
122 """Resolve the requested source set: explicit --search wins, then the
123 LAST30DAYS_DEFAULT_SEARCH config key (env var or .env file), then None
124 (per-query default behavior). The config fallback lets users pin a fixed
125 source set that survives upgrades without patching SKILL.md (#442).
126 """
127 if args_search:
128 return parse_search_flag(args_search)
129 default_search = (config.get("LAST30DAYS_DEFAULT_SEARCH") or "").strip()
130 if default_search:
131 return parse_search_flag(default_search, flag_name="LAST30DAYS_DEFAULT_SEARCH")
132 return None
133
134
135 def add_deep_research_source(
136 requested_sources: list[str] | None,
137 ) -> list[str] | None:
138 """Add Perplexity without replacing the default-source sentinel.
139
140 ``None`` means that the planner can use the normal configured source set.
141 Deep Research enables Perplexity through ``INCLUDE_SOURCES`` separately, so
142 converting this sentinel to ``["perplexity"]`` would suppress every normal
143 source.
144 """
145 if requested_sources is None:
146 return None
147 if "perplexity" in requested_sources:
148 return requested_sources
149 return [*requested_sources, "perplexity"]
150
151
152 def enable_deep_research_source(config: dict) -> None:
153 """Enable the exact Perplexity token or reject a hard exclusion."""
154 excluded = {
155 token.strip().lower()
156 for token in str(config.get("EXCLUDE_SOURCES") or "").split(",")
157 if token.strip()
158 }
159 if "perplexity" in excluded:
160 raise ValueError(
161 "--deep-research conflicts with EXCLUDE_SOURCES=perplexity"
162 )
163
164 include = str(config.get("INCLUDE_SOURCES") or "")
165 tokens = [token.strip() for token in include.split(",") if token.strip()]
166 if "perplexity" not in {token.lower() for token in tokens}:
167 tokens.append("perplexity")
168 config["INCLUDE_SOURCES"] = ",".join(tokens)
169
170
171 def plan_has_explicit_trustpilot_domain(comp_plan: dict | None) -> bool:
172 """True when any --competitors-plan entry pins a trustpilot_domain."""
173 if not comp_plan:
174 return False
175 for entry in comp_plan.values():
176 if not isinstance(entry, dict):
177 continue
178 domain = entry.get("trustpilot_domain")
179 if isinstance(domain, str) and domain.strip():
180 return True
181 return False
182
183
184 def activate_trustpilot_for_explicit_domain(
185 config: dict,
186 requested_sources: list[str] | None,
187 *,
188 reason: str,
189 ) -> list[str] | None:
190 """Activate the opt-in Trustpilot source when the user pinned a domain.
191
192 Passing ``--trustpilot-domain`` (or a plan-level ``trustpilot_domain``) is
193 unambiguous intent — silently ignoring it when Trustpilot is not in
194 ``INCLUDE_SOURCES`` / ``--search`` is the #873 failure mode. Auto-resolve
195 hints must not call this helper.
196
197 ``EXCLUDE_SOURCES=trustpilot`` still wins. Mutates ``config`` in place and
198 returns the (possibly extended) ``requested_sources`` list.
199 """
200 excluded = {
201 token.strip().lower()
202 for token in str(config.get("EXCLUDE_SOURCES") or "").split(",")
203 if token.strip()
204 }
205 if "trustpilot" in excluded:
206 sys.stderr.write(
207 f"[Trustpilot] {reason} ignored: trustpilot is in EXCLUDE_SOURCES\n"
208 )
209 return requested_sources
210
211 include = str(config.get("INCLUDE_SOURCES") or "")
212 tokens = [token.strip() for token in include.split(",") if token.strip()]
213 if "trustpilot" not in {token.lower() for token in tokens}:
214 tokens.append("trustpilot")
215 config["INCLUDE_SOURCES"] = ",".join(tokens)
216 sys.stderr.write(
217 f"[Trustpilot] {reason} activated trustpilot source "
218 "(add to INCLUDE_SOURCES permanently to skip this auto-enable)\n"
219 )
220
221 if requested_sources is not None and "trustpilot" not in requested_sources:
222 requested_sources = [*requested_sources, "trustpilot"]
223 return requested_sources
224
225
226 def activate_telegram_for_explicit_sources(
227 config: dict,
228 requested_sources: list[str] | None,
229 *,
230 channels: str,
231 ) -> list[str] | None:
232 """Activate the opt-in Telegram source when the user pinned channel(s).
233
234 Passing ``--telegram-sources`` is unambiguous intent — silently ignoring it
235 when Telegram is not in ``INCLUDE_SOURCES`` / ``--search`` is the same
236 failure mode as #873 (Trustpilot). Auto-activate the source.
237
238 ``EXCLUDE_SOURCES=telegram`` still wins. Mutates ``config`` in place and
239 returns the (possibly extended) ``requested_sources`` list.
240 """
241 excluded = {
242 token.strip().lower()
243 for token in str(config.get("EXCLUDE_SOURCES") or "").split(",")
244 if token.strip()
245 }
246 if "telegram" in excluded:
247 sys.stderr.write(
248 f"[Telegram] --telegram-sources={channels} ignored: telegram is in EXCLUDE_SOURCES\n"
249 )
250 return requested_sources
251
252 config["TELEGRAM_SOURCES"] = channels
253
254 include = str(config.get("INCLUDE_SOURCES") or "")
255 tokens = [token.strip() for token in include.split(",") if token.strip()]
256 if "telegram" not in {token.lower() for token in tokens}:
257 tokens.append("telegram")
258 config["INCLUDE_SOURCES"] = ",".join(tokens)
259 sys.stderr.write(
260 f"[Telegram] --telegram-sources={channels} activated telegram source "
261 "(add to INCLUDE_SOURCES permanently to skip this auto-enable)\n"
262 )
263
264 if requested_sources is not None and "telegram" not in requested_sources:
265 requested_sources = [*requested_sources, "telegram"]
266 return requested_sources
267
268
269 def slugify(value: str, max_length: int = 180) -> str:
270 slug = re.sub(r"[^a-z0-9]+", "-", value.lower()).strip("-")
271 if len(slug) > max_length:
272 # Filenames built from long topics can exceed the OS 255-byte limit
273 # (macOS errno 63). Truncate and append a hash of the full value so
274 # distinct long topics still get distinct, deterministic names.
275 digest = hashlib.sha1(slug.encode("utf-8")).hexdigest()[:10]
276 slug = f"{slug[:max_length].rstrip('-')}-{digest}"
277 return slug or "last30days"
278
279
280 def sanitize_suffix(suffix: str) -> str:
281 """Sanitize a user-provided ``--save-suffix`` into a path-safe token.
282
283 The suffix is glued directly into the saved-report filename, so restrict it
284 to the same ``[a-z0-9-]`` class as the topic slug. This neutralizes path
285 separators and parent refs (``/``, ``..``) so a suffix can never escape the
286 save directory, while leaving ordinary values ('v3', 'gemini', a client
287 slug) unchanged. Unlike ``slugify`` there is no fallback token: a suffix
288 that sanitizes to nothing simply drops, yielding no suffix part.
289 """
290 return re.sub(r"[^a-z0-9]+", "-", suffix.lower()).strip("-")
291
292
293 def _report_has_private_corpus(report: schema.Report) -> bool:
294 items_by_source = getattr(report, "items_by_source", {})
295 if isinstance(items_by_source, dict) and items_by_source.get("corpus"):
296 return True
297 candidates = getattr(report, "ranked_candidates", ())
298 if not isinstance(candidates, (list, tuple)):
299 return False
300 return any(
301 candidate.source == "corpus"
302 or any(item.source == "corpus" for item in candidate.source_items)
303 for candidate in candidates
304 )
305
306
307 def _ensure_output_directory(path: Path, *, private: bool) -> None:
308 if not private:
309 path.mkdir(parents=True, exist_ok=True)
310 return
311 missing: list[Path] = []
312 current = path
313 while not current.exists():
314 missing.append(current)
315 current = current.parent
316 path.mkdir(parents=True, exist_ok=True, mode=0o700)
317 for directory in missing:
318 directory.chmod(0o700)
319
320
321 def save_output(
322 report: schema.Report,
323 emit: str,
324 save_dir: str,
325 suffix: str = "",
326 synthesis_md: str | None = None,
327 topic_override: str | None = None,
328 rendered_content: str | None = None,
329 json_profile: str = "agent",
330 register: str = "default",
331 private: bool | None = None,
332 render_fn: Callable[[Path], str] | None = None,
333 ) -> Path:
334 from datetime import datetime
335 path = Path(save_dir).expanduser().resolve()
336 slug = slugify(topic_override or report.topic)
337 extension = "json" if emit == "json" else "html" if emit == "html" else "md"
338 raw_label = "raw-html" if emit == "html" else "raw"
339 safe_suffix = sanitize_suffix(suffix)
340 suffix_part = f"-{safe_suffix}" if safe_suffix else ""
341 base = path / f"{slug}-{raw_label}{suffix_part}.{extension}"
342 date_str = datetime.now().strftime('%Y-%m-%d')
343 candidates = [base]
344 candidates.append(path / f"{slug}-{raw_label}{suffix_part}-{date_str}.{extension}")
345 for i in range(1, 100):
346 candidates.append(path / f"{slug}-{raw_label}{suffix_part}-{date_str}-{i}.{extension}")
347 # Markdown saves keep the complete debug artifact. JSON and HTML preserve
348 # their requested wire format so file extensions match their content.
349 # When render_fn is supplied, content is produced after O_EXCL allocates
350 # the candidate. This lets the footer cite the file actually written
351 # without racing a separate filesystem probe.
352 if render_fn is None:
353 if rendered_content is not None:
354 static_content = rendered_content
355 elif emit in {"json", "html"}:
356 static_content = emit_output(
357 report,
358 emit,
359 synthesis_md=synthesis_md,
360 json_profile=json_profile,
361 register=register,
362 )
363 else:
364 static_content = render.render_full(report)
365 private_corpus = _report_has_private_corpus(report) or bool(private)
366 _ensure_output_directory(path, private=private_corpus)
367 for candidate in candidates:
368 try:
369 fd = os.open(
370 candidate,
371 os.O_CREAT | os.O_EXCL | os.O_WRONLY,
372 0o600 if private_corpus else 0o644,
373 )
374 except FileExistsError:
375 continue
376 try:
377 with os.fdopen(fd, "wb") as f:
378 content = render_fn(candidate) if render_fn is not None else static_content
379 f.write(content.encode("utf-8"))
380 except BaseException:
381 # Deferred rendering happens after the candidate is reserved. Do
382 # not leave an empty or partial report if rendering or writing fails.
383 try:
384 candidate.unlink(missing_ok=True)
385 except OSError:
386 pass
387 raise
388 if candidate.suffix.lower() == ".md":
389 try:
390 from lib import library, library_index
391
392 save_root = candidate.parent.resolve()
393 if save_root == Path(library.DEFAULT_MEMORY_DIR).expanduser().resolve():
394 library_index.sync_library(save_root)
395 else:
396 # A scoped save must sync a per-directory index with the
397 # same paths scoped search uses; syncing the shared DB
398 # from one scope's scan would prune other scopes' rows.
399 library_index.sync_library(
400 save_root,
401 save_root / "briefings",
402 db_path=save_root / ".last30days-library.db",
403 )
404 except (library_index.LibrarySearchUnavailable, OSError, sqlite3.DatabaseError):
405 # Saving research must not depend on the optional local index;
406 # `library search` reports a clear capability error on demand.
407 pass
408 return candidate
409 # Fallback: all 101 candidates existed (extremely unlikely).
410 raise RuntimeError(
411 f"save_output: could not find a unique filename after 101 attempts in {path}"
412 )
413
414
415 def save_rendered_output(
416 rendered_content: str,
417 output_file: str,
418 *,
419 private: bool = False,
420 ) -> Path:
421 out_path = Path(output_file).expanduser().resolve()
422 _ensure_output_directory(out_path.parent, private=private)
423 if private and out_path.exists():
424 out_path.chmod(0o600)
425 fd = os.open(
426 out_path,
427 os.O_CREAT | os.O_TRUNC | os.O_WRONLY,
428 0o600 if private else 0o644,
429 )
430 with os.fdopen(fd, "w", encoding="utf-8") as handle:
431 handle.write(rendered_content)
432 if private:
433 out_path.chmod(0o600)
434 return out_path
435
436
437 def _publish_metadata_path(html_path: Path) -> Path:
438 return html_path.with_name(f"{html_path.name}.publish.json")
439
440
441 def _write_publish_metadata(html_path: Path, publish_result: dict[str, object]) -> None:
442 payload = {
443 "url": publish_result.get("url"),
444 "site_id": publish_result.get("site_id"),
445 "status": publish_result.get("status"),
446 "published_at": datetime.datetime.now(datetime.timezone.utc).isoformat(),
447 }
448 _publish_metadata_path(html_path).write_text(json.dumps(payload, indent=2), encoding="utf-8")
449
450
451 def publish_rendered_html(
452 rendered: str,
453 *,
454 password: str | None = None,
455 companion_paths: list[Path] | None = None,
456 ) -> dict[str, object]:
457 from lib import html_publish
458
459 result = html_publish.publish_html(rendered, password=password)
460 metadata_errors: list[str] = []
461 for path in companion_paths or []:
462 try:
463 _write_publish_metadata(path, result)
464 except OSError as exc:
465 metadata_errors.append(f"{path}: {exc}")
466 if metadata_errors:
467 result = dict(result)
468 result["_metadata_errors"] = metadata_errors
469 return result
470
471
472 def _publish_password_for_args(
473 args: argparse.Namespace,
474 config: dict[str, object] | None = None,
475 ) -> str | None:
476 return (
477 args.publish_password
478 or env.read_secret_env("LAST30DAYS_PUBLISH_PASSWORD")
479 or (config or {}).get("LAST30DAYS_PUBLISH_PASSWORD")
480 or None
481 )
482
483
484 def emit_output(
485 report: schema.Report,
486 emit: str,
487 fun_level: str = "medium",
488 save_path: str | None = None,
489 synthesis_md: str | None = None,
490 json_profile: str = "agent",
491 register: str = "default",
492 ) -> str:
493 if emit == "json":
494 payload = (
495 schema.to_dict(report)
496 if json_profile == "raw"
497 else schema.to_agent_export(report)
498 )
499 return json.dumps(payload, indent=2, sort_keys=True)
500 if emit == "html":
501 return html_render.render_html(
502 report,
503 fun_level=fun_level,
504 save_path=save_path,
505 synthesis_md=synthesis_md,
506 register=register,
507 )
508 if emit in {"compact", "md"}:
509 return render.render_compact(
510 report,
511 fun_level=fun_level,
512 save_path=save_path,
513 register=register,
514 )
515 if emit == "context":
516 return render.render_context(report)
517 if emit == "brief":
518 return render.render_brief(report)
519 raise SystemExit(f"Unsupported emit mode: {emit}")
520
521
522 def emit_comparison_output(
523 entity_reports: list[tuple[str, schema.Report]],
524 emit: str,
525 fun_level: str = "medium",
526 save_path: str | None = None,
527 synthesis_md: str | None = None,
528 json_profile: str = "agent",
529 ) -> str:
530 if emit == "json":
531 payload = {
532 "comparison": True,
533 "entities": [label for label, _ in entity_reports],
534 "reports": [
535 {
536 "entity": label,
537 "report": (
538 schema.to_dict(report)
539 if json_profile == "raw"
540 else schema.to_agent_export(report)
541 ),
542 }
543 for label, report in entity_reports
544 ],
545 }
546 if json_profile == "agent":
547 payload["schema_version"] = schema.AGENT_EXPORT_SCHEMA_VERSION
548 return json.dumps(payload, indent=2, sort_keys=True)
549 if emit == "html":
550 return html_render.render_html_comparison(
551 entity_reports,
552 fun_level=fun_level,
553 save_path=save_path,
554 synthesis_md=synthesis_md,
555 )
556 if emit in {"compact", "md"}:
557 return render.render_comparison_multi(
558 entity_reports, fun_level=fun_level, save_path=save_path,
559 )
560 if emit == "context":
561 return render.render_comparison_multi_context(entity_reports)
562 raise SystemExit(f"Unsupported emit mode: {emit}")
563
564
565 def comparison_topic(entity_reports: list[tuple[str, schema.Report]]) -> str:
566 return " vs ".join(label for label, _ in entity_reports)
567
568
569 def comparison_label_key(label: str) -> str:
570 """Normalize an entity label for duplicate detection.
571
572 Comparison labels double as keys in the fan-out's results dict, so two
573 entities differing only in case, surrounding space, or a repeated space
574 collide there while still looking distinct on the command line. Spaces
575 are collapsed, never stripped: "Open AI" stays distinct from "OpenAI".
576 """
577 return " ".join(label.split()).casefold()
578
579
580 def compute_save_path_display(save_dir: str, topic: str, suffix: str, emit: str) -> str:
581 """Compute the user-friendly save path string that will be shown in the footer.
582
583 Uses ~ when the saved file is under the user's home directory; otherwise
584 returns the absolute path.
585 """
586 from pathlib import Path as _Path
587 path = _Path(save_dir).expanduser().resolve()
588 slug = slugify(topic)
589 extension = "json" if emit == "json" else "html" if emit == "html" else "md"
590 raw_label = "raw-html" if emit == "html" else "raw"
591 safe_suffix = sanitize_suffix(suffix)
592 suffix_part = f"-{safe_suffix}" if safe_suffix else ""
593 raw = path / f"{slug}-{raw_label}{suffix_part}.{extension}"
594 try:
595 home = _Path.home().resolve()
596 relative = raw.relative_to(home)
597 return f"~/{relative.as_posix()}"
598 except ValueError:
599 return raw.as_posix()
600
601
602 def compute_output_path_display(output_file: str) -> str:
603 """Compute the user-friendly explicit output path shown in render footers."""
604 raw = Path(output_file).expanduser().resolve()
605 try:
606 home = Path.home().resolve()
607 relative = raw.relative_to(home)
608 return f"~/{relative.as_posix()}"
609 except ValueError:
610 return raw.as_posix()
611
612
613 def read_synthesis_file(path: str) -> str:
614 try:
615 return Path(path).expanduser().read_text(encoding="utf-8")
616 except OSError as exc:
617 sys.stderr.write(f"[last30days] Cannot read --synthesis-file: {exc}\n")
618 raise SystemExit(2)
619
620
621 def _scoped_store_db(args: argparse.Namespace) -> Path | None:
622 """Scoped runs write findings inside the save dir, matching scoped reads."""
623 save_dir = getattr(args, "save_dir", None)
624 if save_dir:
625 return Path(save_dir).expanduser().resolve() / "research.db"
626 return None
627
628
629 def persist_report(report: schema.Report, store_db: Path | None = None) -> dict[str, int]:
630 import store
631
632 private_corpus = _report_has_private_corpus(report)
633 with store.scoped_db(store_db):
634 if private_corpus:
635 store.ensure_private_db_files()
636 store.init_db()
637 if private_corpus:
638 store.ensure_private_db_files()
639 topic_row = store.add_topic(report.topic, update_existing=False)
640 topic_id = topic_row["id"]
641 source_mode = ",".join(sorted(report.items_by_source)) or "v3"
642 run_id = store.record_run(topic_id, source_mode=source_mode, status="running")
643 try:
644 findings = store.findings_from_report(report)
645 if private_corpus:
646 store.ensure_private_db_files()
647 counts = store.store_findings(run_id, topic_id, findings)
648 store.update_run(
649 run_id,
650 status="completed",
651 findings_new=counts["new"],
652 findings_updated=counts["updated"],
653 )
654 return counts
655 except Exception as exc:
656 store.update_run(run_id, status="failed", error_message=str(exc)[:500])
657 raise
658 finally:
659 if private_corpus:
660 store.ensure_private_db_files()
661
662
663 def build_parser() -> argparse.ArgumentParser:
664 parser = argparse.ArgumentParser(
665 description="Research a topic across live social, market, and grounded web sources.",
666 allow_abbrev=False,
667 )
668 parser.add_argument("topic", nargs="*", help="Research topic")
669 parser.add_argument("--emit", default="compact", choices=["compact", "json", "context", "md", "html", "brief"])
670 parser.add_argument(
671 "--register",
672 choices=registers.REGISTER_NAMES,
673 default=None,
674 help="Audience synthesis preset for the standard brief (default, exec, dev, creator, eli5)",
675 )
676 parser.add_argument(
677 "--json-profile",
678 default="agent",
679 choices=["agent", "raw"],
680 help="JSON export profile for --emit=json (default: agent)",
681 )
682 parser.add_argument("--search", help="Comma-separated source list")
683 parser.add_argument("--quick", action="store_true", help="Lower-latency retrieval profile")
684 parser.add_argument("--deep", action="store_true", help="Higher-recall retrieval profile")
685 freshness_group = parser.add_mutually_exclusive_group()
686 freshness_group.add_argument(
687 "--verify-freshness",
688 action="store_true",
689 default=None,
690 help="Re-check source-grounded claims after research, or verify the cached report when no topic is supplied",
691 )
692 freshness_group.add_argument(
693 "--no-verify-freshness",
694 dest="verify_freshness",
695 action="store_false",
696 help="Disable freshness verification configured by LAST30DAYS_VERIFY_FRESHNESS",
697 )
698 parser.add_argument(
699 "--drill",
700 metavar="TARGET",
701 help="Deep follow-up on a cluster from the fresh last-report.json cache",
702 )
703 parser.add_argument(
704 "--discover",
705 metavar="DOMAIN",
706 nargs="?",
707 const="",
708 default=None,
709 help=(
710 "Sweep river listings and rank the topics accelerating in a domain; "
711 "each survivor gets a full research pass. Bare --discover (no domain) "
712 "runs global trending across every feed's hot list"
713 ),
714 )
715 parser.add_argument(
716 "--discover-shallow",
717 action="store_true",
718 help=(
719 "Skip the per-topic research pass during --discover: rank on listing "
720 "evidence only (faster, thinner; the confidence floor still applies)"
721 ),
722 )
723 parser.add_argument(
724 "--nominate-only",
725 action="store_true",
726 help=(
727 "Leg 1 of the host-judged discovery protocol: sweep, write the "
728 "nominations bundle for host judgment, and stop (no judging, no "
729 "enrichment). Requires --discover"
730 ),
731 )
732 parser.add_argument(
733 "--judgments",
734 metavar="PATH",
735 help=(
736 "Leg 2 of the discovery protocol: resume from the nominations "
737 "bundle, applying the host judgments file at PATH. Requires "
738 "--discover"
739 ),
740 )
741 parser.add_argument(
742 "--finalize",
743 action="store_true",
744 help=(
745 "Leg 3 of the discovery protocol: apply host angles, render the "
746 "final discovery brief, and record the topic queue. Requires "
747 "--discover"
748 ),
749 )
750 parser.add_argument(
751 "--angles",
752 metavar="PATH",
753 help=(
754 "Optional host angles file for --discover --finalize (omitting it "
755 "ships the brief without angle lines)"
756 ),
757 )
758 parser.add_argument("--debug", action="store_true", help="Enable HTTP debug logging")
759 parser.add_argument("--mock", action="store_true", help="Use mock retrieval fixtures")
760 parser.add_argument(
761 "--record-fixtures",
762 metavar="DIR",
763 help=argparse.SUPPRESS,
764 )
765 parser.add_argument("--diagnose", action="store_true", help="Print provider and source availability")
766 parser.add_argument("--preflight", action="store_true",
767 help="Print a safe human-readable permission preflight")
768 parser.add_argument("--welcome", action="store_true",
769 help="Print the first-run welcome text (engine-owned; relay verbatim)")
770 parser.add_argument("--preflight-report-on-save-dir", help=argparse.SUPPRESS)
771 parser.add_argument("--no-browser-cookies", action="store_true",
772 help="Disable browser-cookie extraction even when FROM_BROWSER is configured")
773 parser.add_argument("--save-dir", help="Optional directory for saving the rendered output")
774 parser.add_argument(
775 "--resolve-save-dir", action="store_true",
776 help="Print the skill save directory from flags/config, then exit without research",
777 )
778 parser.add_argument(
779 "--corpus",
780 action="append",
781 default=[],
782 metavar="DIR",
783 help="Add a local .md/.txt/.pdf directory as a private ranked source (repeatable)",
784 )
785 parser.add_argument(
786 "--corpus-all-time",
787 action="store_true",
788 help="Include matching corpus files older than the research window",
789 )
790 parser.add_argument("--output", help="Optional exact file path for saving the rendered output")
791 parser.add_argument("--synthesis-file", help="Markdown synthesis to embed in --emit=html output")
792 parser.add_argument("--publish-html", action="store_true",
793 help="Publish --emit=html output to ht-ml.app (explicit opt-in; public by default)")
794 parser.add_argument("--publish", action="store_true",
795 help="With 'library feed', publish the HTML index and briefs (explicit opt-in; public by default); feed.xml remains local")
796 parser.add_argument("--publish-password",
797 help="Optional shared password for --publish-html or 'library feed --publish'; prefer LAST30DAYS_PUBLISH_PASSWORD to avoid exposing secrets in process lists")
798 parser.add_argument("--store", action="store_true", help="Persist ranked findings to the SQLite research store")
799 parser.add_argument("--x-handle", help="X handle for targeted supplemental search")
800 parser.add_argument("--x-related", help="Comma-separated related X handles (searched with lower weight)")
801 parser.add_argument(
802 "--x-posts",
803 dest="x_posts",
804 metavar="PATH",
805 help=(
806 "Path to a last30days-x-posts/1 JSON envelope of posts the hosting "
807 "model fetched through its X connector; replaces the engine's X "
808 "fetch for this run. A file path only (never inline JSON); on a "
809 "comparison run use the per-entity x_posts field of --competitors-plan."
810 ),
811 )
812 parser.add_argument("--web-backend", default="auto",
813 choices=["auto", "brave", "exa", "serper", "parallel", "parallel-mcp", "keyless", "none"],
814 help="Web search backend (default: auto; parallel-mcp explicitly opts into the "
815 "anonymous hosted MCP; keyless forces the zero-key floor)")
816 parser.add_argument("--perplexity-search-type", choices=["web", "fast"],
817 help="Search backend for direct Perplexity Search API and Agent web_search; overrides LAST30DAYS_PERPLEXITY_SEARCH_TYPE. Does not enable the paid source or select an Agent preset.")
818 parser.add_argument("--deep-research", action="store_true",
819 help="Use at most one Perplexity Deep Research run. Direct PERPLEXITY_API_KEY uses the Agent API background path; OPENROUTER_API_KEY keeps the synchronous Sonar fallback; cannot be combined with competitor or vs-mode.")
820 parser.add_argument("--hiring-signals", action="store_true",
821 help="Analyze public jobs/careers postings as evidence-backed company focus signals.")
822 parser.add_argument("--plan", help="JSON query plan (skips internal LLM planner). Can be a JSON string or a file path.")
823 parser.add_argument("--save-suffix", help="Suffix for saved output filename (e.g., 'gemini' → kanye-west-raw-gemini.md)")
824 parser.add_argument("--subreddits", help="Comma-separated broad/category subreddit names to search (e.g., SaaS,Entrepreneur)")
825 parser.add_argument("--dedicated-subreddits", help="Comma-separated entity-home subreddit names (e.g., Kanye,WestSubEver). Pulled in full (top+hot+new) and exempt from the relevance floor since the whole sub is the topic.")
826 parser.add_argument("--tiktok-hashtags", help="Comma-separated TikTok hashtags without # (e.g., tella,screenrecording)")
827 parser.add_argument("--tiktok-creators", help="Comma-separated TikTok creator handles (e.g., TellaHQ,taborplace)")
828 parser.add_argument("--ig-creators", help="Comma-separated Instagram creator handles (e.g., tella.tv,laborstories)")
829 parser.add_argument(
830 "--days",
831 "--lookback-days",
832 dest="lookback_days",
833 type=int,
834 default=None,
835 help="Number of days to look back for research (default: 30, watchlist uses 90)",
836 )
837 parser.add_argument(
838 "--as-of",
839 dest="as_of_date",
840 type=parse_as_of_date_arg,
841 help=(
842 "End date for the lookback window in YYYY-MM-DD format. "
843 "When set, --days looks back from this date instead of today."
844 ),
845 )
846 parser.add_argument("--max-results", dest="max_results", type=int,
847 help="Override the final ranked-pool cap (pool_limit/rerank_limit) from the depth profile. "
848 "Use for high-volume topics where the default (deep=60) under-covers. See issue #716.")
849 parser.add_argument("--max-per-source", dest="max_per_source", type=int,
850 help="Override the per-stream cap (per_stream_limit) applied to each (source, subquery) before "
851 "pooling. Raising it increases unique-item yield when one source has many relevant items. "
852 "See issue #716.")
853 parser.add_argument("--max-source-fetches", dest="max_source_fetches", type=int,
854 help="Override the per-source fetch cap (MAX_SOURCE_FETCHES, default x=2) that limits how many "
855 "subqueries actually fetch a capped source. Raise it so every X subquery in a multi-angle "
856 "--plan runs instead of just the first two. See issue #716.")
857 parser.add_argument("--auto-resolve", action="store_true",
858 help="Use web search to discover subreddits/handles before planning (for platforms without WebSearch)")
859 parser.add_argument("--github-user", help="GitHub username for person-mode search (e.g., steipete)")
860 parser.add_argument("--github-repo", help="Comma-separated owner/repo for project-mode search (e.g., openclaw/openclaw,paperclipai/paperclip)")
861 parser.add_argument(
862 "--trustpilot-domain",
863 help=(
864 "Trustpilot review-page domain for the topic (e.g., www.thriftbooks.com). "
865 "Used verbatim, bypasses the brand-shape gate, and auto-activates the "
866 "opt-in Trustpilot source for this run (unless EXCLUDE_SOURCES=trustpilot). "
867 "Find the domain with `trustpilot-pp-cli search '<name>'`."
868 ),
869 )
870 parser.add_argument(
871 "--amazon-query",
872 help=(
873 "Product keyword the amazon source searches, when that source is active. "
874 "Defaults to the topic. Supply it whenever the topic is not the product: "
875 "a person topic searches their company's product line "
876 "(--amazon-query='June Oven'), and a brand searches brand-plus-category "
877 "(--amazon-query='Weber grill', not 'Weber' -- a bare brand keyword lands "
878 "on an ad-heavy page that can miss the brand's own bestsellers). "
879 "Requires the brightdata CLI on PATH and logged in."
880 ),
881 )
882 parser.add_argument(
883 "--meta-ads-page",
884 help=(
885 "Meta Ad Library page id for the topic's advertiser, when the meta_ads "
886 "source is active. Skips name-based page resolution and its discovery "
887 "credit. Accepts a bare numeric page id (e.g. 123456789012345) or an Ad "
888 "Library URL carrying view_all_page_id. A facebook.com vanity URL is not "
889 "a page id and is rejected. Use it when a brand advertises under product "
890 "names, or when resolution picked the wrong company."
891 ),
892 )
893 parser.add_argument(
894 "--telegram-sources",
895 help=(
896 "Comma-separated list of public Telegram channel handles or t.me URLs. "
897 "Auto-activates the opt-in Telegram source for this run. "
898 "Accepts: bare handle (aipost), @handle (@aipost), "
899 "t.me URL (https://t.me/aipost), or preview URL (https://t.me/s/aipost). "
900 "Rejects joinchat links and numeric -100 supergroup IDs."
901 ),
902 )
903 parser.add_argument(
904 "--competitors",
905 nargs="?",
906 const=2,
907 type=int,
908 default=None,
909 metavar="N",
910 help="Auto-discover N competitor entities and fan out last30days across all of them as a comparison (default N=2 → 3-way: original + 2 peers; range 1..6). Use --competitors-list to override discovery.",
911 )
912 parser.add_argument(
913 "--competitors-list",
914 dest="competitors_list",
915 help="Comma-separated competitor entities to skip discovery (e.g., 'Anthropic,xAI,Google Gemini'). Implies --competitors.",
916 )
917 parser.add_argument(
918 "--polymarket-keywords",
919 dest="polymarket_keywords",
920 help=(
921 "Comma-separated keywords that Polymarket market titles must match "
922 "to be included. Use for ambiguous single-token topics like 'Warriors' "
923 "(nba,gsw,golden-state) to filter out Glasgow Warriors rugby, Honor "
924 "of Kings Rogue Warriors, etc. When omitted, Polymarket returns all "
925 "matching markets — so expect cross-entity noise on generic topics."
926 ),
927 )
928 parser.add_argument(
929 "--competitors-plan",
930 dest="competitors_plan",
931 help=(
932 "JSON mapping of per-entity Step 0.55 targeting for competitor / vs-mode "
933 "sub-runs. Schema: {entity_name: {x_handle?, x_related?, subreddits?, "
934 "github_user?, github_repos?, context?}}. Accepts inline JSON or a file "
935 "path. Implies --competitors. Preferred over --competitors-list when the "
936 "hosting model has already resolved per-entity handles and subs."
937 ),
938 )
939 return parser
940
941
942 def parse_competitors_plan(raw: str | None) -> dict[str, dict]:
943 """Parse a --competitors-plan argument into a {entity_name_lower: plan_entry} dict.
944
945 Accepts inline JSON or a file path (matches --plan). Returns {} on None/empty.
946 Validation: top-level must be a dict; each value must be a dict. Unknown fields
947 in entry values log a warning but do not abort. Invalid JSON or non-dict shape
948 raises SystemExit(2) with a clear stderr message.
949 """
950 if not raw:
951 return {}
952 plan_str = raw
953 if os.path.isfile(plan_str):
954 try:
955 with open(plan_str, encoding="utf-8") as f:
956 plan_str = f.read()
957 except (OSError, UnicodeDecodeError) as exc:
958 sys.stderr.write(f"[CompetitorsPlan] Cannot read plan file: {exc}\n")
959 raise SystemExit(2)
960 try:
961 parsed = json.loads(plan_str)
962 except json.JSONDecodeError as exc:
963 sys.stderr.write(f"[CompetitorsPlan] Invalid JSON: {exc}\n")
964 raise SystemExit(2)
965 if not isinstance(parsed, dict):
966 sys.stderr.write(
967 f"[CompetitorsPlan] Top-level must be a dict of "
968 f"{{entity: {{targeting}}}}, got {type(parsed).__name__}\n"
969 )
970 raise SystemExit(2)
971 known_fields = {
972 "x_handle", "x_related", "subreddits",
973 "github_user", "github_repos", "trustpilot_domain", "context",
974 "x_posts",
975 }
976 normalized: dict[str, dict] = {}
977 for entity, entry in parsed.items():
978 if not isinstance(entry, dict):
979 sys.stderr.write(
980 f"[CompetitorsPlan] Entry for {entity!r} must be a dict, "
981 f"got {type(entry).__name__}; skipping.\n"
982 )
983 continue
984 unknown = set(entry.keys()) - known_fields
985 if unknown:
986 sys.stderr.write(
987 f"[CompetitorsPlan] Unknown fields in {entity!r}: "
988 f"{sorted(unknown)}; ignoring.\n"
989 )
990 normalized[entity.strip().lower()] = {
991 **{k: v for k, v in entry.items() if k in known_fields},
992 "_name": entity.strip(),
993 }
994 return normalized
995
996
997 def subrun_kwargs_for(
998 entity: str,
999 plan_entry: dict,
1000 *,
1001 resolved: dict,
1002 ) -> dict:
1003 """Build an explicit per-entity kwargs dict for pipeline.run().
1004
1005 Plan values win over auto_resolve values. Returns keys for all per-entity
1006 targeting flags so callers never fall through to closure defaults.
1007
1008 This helper is the single source of truth for sub-run kwargs — main-topic
1009 flags can only leak if a caller bypasses it.
1010 """
1011 def _choose(plan_key: str, resolved_key: str | None = None):
1012 if plan_key in plan_entry and plan_entry[plan_key]:
1013 return plan_entry[plan_key]
1014 if resolved_key is not None and resolved.get(resolved_key):
1015 return resolved[resolved_key]
1016 return None
1017
1018 x_handle = _choose("x_handle", "x_handle")
1019 if isinstance(x_handle, str):
1020 x_handle = x_handle.lstrip("@") or None
1021
1022 subreddits = _choose("subreddits", "subreddits")
1023 if isinstance(subreddits, list):
1024 subreddits = [s.strip().removeprefix("r/") for s in subreddits if s.strip()] or None
1025
1026 x_related = plan_entry.get("x_related")
1027 if isinstance(x_related, list):
1028 x_related = [h.strip().lstrip("@") for h in x_related if h.strip()] or None
1029 else:
1030 x_related = None
1031
1032 github_user = _choose("github_user", "github_user")
1033 if isinstance(github_user, str):
1034 github_user = github_user.lstrip("@").lower() or None
1035
1036 github_repos = _choose("github_repos", "github_repos")
1037 if isinstance(github_repos, list):
1038 github_repos = [r.strip() for r in github_repos if r.strip() and "/" in r.strip()] or None
1039
1040 trustpilot_domain = _choose("trustpilot_domain", "trustpilot_domain")
1041 if isinstance(trustpilot_domain, str):
1042 trustpilot_domain = trustpilot_domain.strip() or None
1043 # Provenance: a plan-supplied domain is user-set (verbatim-final); one that
1044 # only came from auto_resolve is a hint that retries via search on a miss.
1045 trustpilot_domain_is_hint = bool(
1046 trustpilot_domain and not plan_entry.get("trustpilot_domain")
1047 )
1048
1049 context = plan_entry.get("context") or resolved.get("context") or ""
1050
1051 return {
1052 "x_handle": x_handle,
1053 "x_related": x_related,
1054 "subreddits": subreddits,
1055 "github_user": github_user,
1056 "github_repos": github_repos,
1057 "trustpilot_domain": trustpilot_domain,
1058 "_trustpilot_domain_is_hint": trustpilot_domain_is_hint,
1059 "_context": context,
1060 }
1061
1062
1063 COMPETITORS_MIN = competitors_mod.COMPETITORS_MIN
1064 COMPETITORS_MAX = competitors_mod.COMPETITORS_MAX
1065 COMPETITORS_DEFAULT = competitors_mod.COMPETITORS_DEFAULT
1066
1067
1068 def truncate_comparison_entities(entities: list[str], *, warn: bool = True) -> list[str]:
1069 """Cap a vs-entity list at COMPARISON_ENTITY_MAX; optionally warn on stderr."""
1070 ceiling = competitors_mod.COMPARISON_ENTITY_MAX
1071 if len(entities) <= ceiling:
1072 return list(entities)
1073 kept = entities[:ceiling]
1074 dropped = entities[ceiling:]
1075 if warn:
1076 sys.stderr.write(
1077 f"[Competitors] vs-topic has {len(entities)} entities; "
1078 f"using first {ceiling}, dropped: {', '.join(dropped)}\n"
1079 )
1080 return kept
1081
1082
1083 def apply_vs_competitor_routing(
1084 topic: str,
1085 *,
1086 competitors_flag: int | None,
1087 comp_enabled: bool,
1088 comp_count: int,
1089 comp_explicit: list[str],
1090 comp_plan: dict[str, dict] | None = None,
1091 ) -> tuple[str, bool, int, list[str]]:
1092 """Apply vs-string / plan routing on top of resolve_competitors_args.
1093
1094 Precedence for *who* runs:
1095 1. ``--competitors-list`` (explicit peers; topic unchanged)
1096 2. Pure discover-N (``--competitors`` without list or plan) — topic
1097 unchanged, even if it contains ``vs``
1098 3. vs-string split (first entity becomes main topic) — used for bare
1099 vs-topics and vs-topic + ``--competitors-plan``
1100 4. ``--competitors-plan`` keys as peers when there is no vs-string
1101 (including when ``--competitors N`` is also set)
1102 """
1103 from lib import planner as _planner
1104
1105 if comp_explicit:
1106 return topic, True, len(comp_explicit), list(comp_explicit)
1107
1108 # Preserve discover-N semantics: numeric flag without plan/list must not
1109 # rewrite a vs-string into named peers.
1110 if competitors_flag is not None and not comp_plan:
1111 return topic, True, comp_count, []
1112
1113 vs_entities = truncate_comparison_entities(
1114 _planner._comparison_entities(topic, uncapped=True),
1115 warn=True,
1116 )
1117 if len(vs_entities) >= 2:
1118 main, peers = vs_entities[0], vs_entities[1:]
1119 sys.stderr.write(
1120 f"[Competitors] vs-mode: routing to N-pass fanout: "
1121 f"{main} vs {' vs '.join(peers)}\n"
1122 )
1123 return main, True, len(peers), peers
1124
1125 if comp_plan:
1126 plan_peers = [
1127 (entry.get("_name") or key)
1128 for key, entry in comp_plan.items()
1129 ]
1130 plan_peers = [name for name in plan_peers if name]
1131 if len(plan_peers) > COMPETITORS_MAX:
1132 sys.stderr.write(
1133 f"[Competitors] --competitors-plan has {len(plan_peers)} entries, "
1134 f"clamping to {COMPETITORS_MAX}.\n"
1135 )
1136 plan_peers = plan_peers[:COMPETITORS_MAX]
1137 return topic, True, len(plan_peers), plan_peers
1138
1139 return topic, comp_enabled, comp_count, comp_explicit
1140
1141
1142 def resolve_competitors_args(args: argparse.Namespace) -> tuple[bool, int, list[str]]:
1143 """Normalize competitors flags into (enabled, count, explicit_list).
1144
1145 - (False, 0, []) when neither flag, list, nor plan is set.
1146 - An explicit ``--competitors-list`` always wins; count is derived from list length.
1147 - ``--competitors-plan`` alone enables mode with an empty peer list; vs-routing
1148 fills peers from the vs-string or plan keys.
1149 - A numeric count outside [1, 6] is clamped with a stderr warning.
1150 - count <= 0 (explicit) raises SystemExit(2).
1151 """
1152 explicit_list: list[str] = []
1153 list_flag_provided = args.competitors_list is not None
1154 if list_flag_provided:
1155 explicit_list = [
1156 entity.strip()
1157 for entity in args.competitors_list.split(",")
1158 if entity.strip()
1159 ]
1160 if not explicit_list:
1161 sys.stderr.write("[Competitors] --competitors-list is empty.\n")
1162 raise SystemExit(2)
1163
1164 competitors_flag = args.competitors
1165 list_present = bool(explicit_list)
1166 flag_present = competitors_flag is not None
1167 plan_present = bool(getattr(args, "competitors_plan", None))
1168
1169 if not list_present and not flag_present and not plan_present:
1170 return False, 0, []
1171
1172 if list_present:
1173 count = len(explicit_list)
1174 if flag_present and competitors_flag != count:
1175 sys.stderr.write(
1176 f"[Competitors] --competitors={competitors_flag} ignored; using "
1177 f"{count} entries from --competitors-list.\n"
1178 )
1179 if count > COMPETITORS_MAX:
1180 sys.stderr.write(
1181 f"[Competitors] --competitors-list has {count} entries, clamping to {COMPETITORS_MAX}.\n"
1182 )
1183 explicit_list = explicit_list[:COMPETITORS_MAX]
1184 count = COMPETITORS_MAX
1185 return True, count, explicit_list
1186
1187 if flag_present:
1188 count = competitors_flag
1189 if count < COMPETITORS_MIN:
1190 sys.stderr.write(
1191 f"[Competitors] --competitors must be >= {COMPETITORS_MIN} (got {count}).\n"
1192 )
1193 raise SystemExit(2)
1194 if count > COMPETITORS_MAX:
1195 sys.stderr.write(
1196 f"[Competitors] --competitors={count} exceeds max {COMPETITORS_MAX}; clamping.\n"
1197 )
1198 count = COMPETITORS_MAX
1199 return True, count, []
1200
1201 # plan_present alone: enable; peers filled by apply_vs_competitor_routing.
1202 return True, 0, []
1203
1204
1205 def _missing_sources_for_promo(diag: dict[str, object]) -> str | None:
1206 available = set(diag.get("available_sources") or [])
1207 missing = []
1208 if "reddit" not in available:
1209 missing.append("reddit")
1210 # X is optional. A successful run without X must reach the research output
1211 # without an authentication or browser-cookie promo in front of it.
1212 # The web promo nudges toward a paid backend for higher-quality web search.
1213 # Grounding is now available keyless on non-native hosts, so key the promo on
1214 # the absence of a *paid* backend, not on grounding availability. Suppress it
1215 # entirely on native-search hosts, where the model's own search is better and
1216 # setting a paid engine key would be the wrong advice.
1217 if not diag.get("native_web_backend") and not diag.get("native_search"):
1218 missing.append("web")
1219 if not missing:
1220 return None
1221 return missing[0]
1222
1223
1224 def _optional_x_omission_text(
1225 diag: dict[str, object],
1226 requested_sources: list[str] | None,
1227 ) -> str | None:
1228 """Return a non-blocking post-result note for a default run without X.
1229
1230 Explicit ``--search`` runs already define their intended source boundary,
1231 so they do not need an omission note. Doctor/diagnose remains the place for
1232 X setup or repair instructions.
1233 """
1234 if requested_sources is not None:
1235 return None
1236 available = set(diag.get("available_sources") or [])
1237 if "x" in available:
1238 return None
1239 return (
1240 "Optional source omitted: X/Twitter was not enabled; research "
1241 "continued with the available sources."
1242 )
1243
1244
1245 def _show_runtime_ui(
1246 report: schema.Report,
1247 progress: ui.ProgressDisplay,
1248 diag: dict[str, object],
1249 suppress_web_promo: bool = False,
1250 ) -> None:
1251 counts = {source: len(items) for source, items in report.items_by_source.items()}
1252 display_sources = list(
1253 dict.fromkeys(
1254 [
1255 *report.query_plan.source_weights.keys(),
1256 *report.items_by_source.keys(),
1257 *report.errors_by_source.keys(),
1258 ]
1259 )
1260 )
1261 progress.end_processing()
1262 progress.show_complete(
1263 source_counts=counts,
1264 display_sources=display_sources,
1265 )
1266 promo = _missing_sources_for_promo(diag)
1267 # The `web` promo nudges users to set BRAVE_API_KEY / SERPER_API_KEY, which
1268 # is wrong advice when a hosting reasoning model (Claude Code, Codex,
1269 # Hermes, Gemini) is driving — those already have WebSearch and can
1270 # pre-resolve Step 0.55 themselves. Suppress the web promo when a hosting
1271 # model signal is present (--plan or --competitors-plan was passed).
1272 if promo:
1273 if suppress_web_promo and promo == "web":
1274 return
1275 if suppress_web_promo and promo == "both":
1276 # "both" means reddit + web both missing; still nudge reddit but
1277 # skip the web line. show_promo has a per-source variant.
1278 progress.show_promo("reddit", diag=diag)
1279 return
1280 progress.show_promo(promo, diag=diag)
1281
1282
1283 REPORT_CACHE_VERSION = "last30days-report-cache/v1"
1284 DEFAULT_REPORT_CACHE_TTL_SECONDS = 3600
1285
1286
1287 def _last_report_cache_path() -> Path | None:
1288 if env.CONFIG_DIR is None:
1289 return None
1290 return env.CONFIG_DIR / "last-report.json"
1291
1292
1293 def _report_cache_ttl_seconds(config: dict[str, object]) -> int:
1294 raw = os.environ.get("LAST30DAYS_REPORT_CACHE_TTL_SECONDS")
1295 if raw is None:
1296 raw = config.get("LAST30DAYS_REPORT_CACHE_TTL_SECONDS")
1297 if raw is None or raw == "":
1298 return DEFAULT_REPORT_CACHE_TTL_SECONDS
1299 try:
1300 return max(0, int(raw))
1301 except (TypeError, ValueError):
1302 return DEFAULT_REPORT_CACHE_TTL_SECONDS
1303
1304
1305 def _is_report_cache_fresh(timestamp: object, ttl_seconds: int) -> bool:
1306 return env.is_timestamp_fresh(timestamp, ttl_seconds)
1307
1308
1309 def _write_last_run(
1310 topic: str,
1311 report: "schema.Report",
1312 entity_reports: list[tuple[str, schema.Report]] | None = None,
1313 *,
1314 x_envelope_sha256: str | None = None,
1315 ) -> bool:
1316 # ``x_envelope_sha256`` binds the cached report to the --x-posts file it
1317 # was built from; _load_last_report_cache misses on any mismatch.
1318 try:
1319 if env.CONFIG_DIR is None:
1320 return False
1321 target = env.CONFIG_DIR
1322 cached_reports = entity_reports or [(report.topic, report)]
1323 has_private_corpus = any(
1324 cached_report.items_by_source.get("corpus")
1325 for _, cached_report in cached_reports
1326 )
1327 _ensure_output_directory(target, private=has_private_corpus)
1328 counts = {source: len(items) for source, items in report.items_by_source.items()}
1329 payload = {
1330 "topic": topic,
1331 "timestamp": datetime.datetime.now(datetime.timezone.utc).isoformat(),
1332 "sources": counts,
1333 "total": sum(counts.values()),
1334 "report_cache": str(target / "last-report.json"),
1335 "comparison": bool(entity_reports),
1336 }
1337 (target / "last-run.json").write_text(json.dumps(payload, indent=2))
1338 cache_payload = {
1339 "schema": REPORT_CACHE_VERSION,
1340 "topic": topic,
1341 "timestamp": payload["timestamp"],
1342 "comparison": bool(entity_reports),
1343 "x_envelope_sha256": x_envelope_sha256 or None,
1344 "reports": [
1345 {"entity": label, "report": schema.to_dict(cached_report)}
1346 for label, cached_report in cached_reports
1347 ],
1348 }
1349 report_cache_path = target / "last-report.json"
1350 report_cache_path.write_text(json.dumps(cache_payload, indent=2))
1351 if has_private_corpus:
1352 report_cache_path.chmod(0o600)
1353 return True
1354 except Exception as exc:
1355 # Never fatal, but never silent either (#787's lesson): callers that
1356 # promise cache state (drill chaining) branch on the return value.
1357 sys.stderr.write(f"[last30days] warning: could not write run cache: {exc}\n")
1358 return False
1359
1360
1361 def _load_last_report_cache(
1362 topic: str | None,
1363 ttl_seconds: int = DEFAULT_REPORT_CACHE_TTL_SECONDS,
1364 *,
1365 x_envelope_sha256: str | None = None,
1366 ) -> tuple[schema.Report, list[tuple[str, schema.Report]] | None, Path] | None:
1367 cache_path = _last_report_cache_path()
1368 if cache_path is None or not cache_path.exists():
1369 return None
1370 try:
1371 payload = json.loads(cache_path.read_text(encoding="utf-8"))
1372 if not isinstance(payload, dict):
1373 raise TypeError("report cache payload must be a JSON object")
1374 if payload.get("schema") != REPORT_CACHE_VERSION:
1375 return None
1376 if not _is_report_cache_fresh(payload.get("timestamp"), ttl_seconds):
1377 return None
1378 # A report built from a --x-posts envelope is only reusable with the
1379 # same envelope content; a digest on either side that does not match
1380 # the other is a miss.
1381 cached_digest = payload.get("x_envelope_sha256") or None
1382 if (cached_digest or x_envelope_sha256) and cached_digest != x_envelope_sha256:
1383 return None
1384 cached_topic = str(payload.get("topic") or "").strip().lower()
1385 if topic is not None and cached_topic != topic.strip().lower():
1386 return None
1387 reports_payload = payload.get("reports") or []
1388 if not reports_payload:
1389 return None
1390 entity_reports = [
1391 (str(item.get("entity") or ""), schema.report_from_dict(item["report"]))
1392 for item in reports_payload
1393 if isinstance(item, dict) and isinstance(item.get("report"), dict)
1394 ]
1395 if not entity_reports:
1396 return None
1397 if payload.get("comparison"):
1398 if len(entity_reports) < 2:
1399 return None
1400 if len(entity_reports) != len(reports_payload):
1401 return None
1402 return entity_reports[0][1], entity_reports, cache_path
1403 return entity_reports[0][1], None, cache_path
1404 except (OSError, json.JSONDecodeError, KeyError, TypeError, ValueError) as exc:
1405 sys.stderr.write(
1406 f"[last30days] Could not read report cache {cache_path}: "
1407 f"{type(exc).__name__}: {exc}\n"
1408 )
1409 return None
1410
1411
1412 def _config_truthy(value: object) -> bool:
1413 return str(value or "").strip().lower() in {"1", "true", "yes", "on"}
1414
1415
1416 def _freshness_enabled(args: argparse.Namespace, config: dict[str, object]) -> bool:
1417 if args.verify_freshness is not None:
1418 return bool(args.verify_freshness)
1419 return _config_truthy(config.get("LAST30DAYS_VERIFY_FRESHNESS"))
1420
1421
1422 def _update_cached_freshness(
1423 cache_path: Path,
1424 report: schema.Report,
1425 entity_reports: list[tuple[str, schema.Report]] | None,
1426 ) -> bool:
1427 """Rewrite cached report bodies without extending the research-cache TTL."""
1428 try:
1429 payload = json.loads(cache_path.read_text(encoding="utf-8"))
1430 if not isinstance(payload, dict) or payload.get("schema") != REPORT_CACHE_VERSION:
1431 return False
1432 existing = payload.get("reports") or []
1433 if entity_reports:
1434 cached_reports = entity_reports
1435 else:
1436 label = (
1437 str(existing[0].get("entity") or report.topic)
1438 if existing and isinstance(existing[0], dict)
1439 else report.topic
1440 )
1441 cached_reports = [(label, report)]
1442 payload["reports"] = [
1443 {"entity": label, "report": schema.to_dict(cached_report)}
1444 for label, cached_report in cached_reports
1445 ]
1446 cache_path.write_text(json.dumps(payload, indent=2), encoding="utf-8")
1447 return True
1448 except (OSError, json.JSONDecodeError, TypeError, ValueError) as exc:
1449 sys.stderr.write(
1450 f"[last30days] warning: could not update freshness cache: {exc}\n"
1451 )
1452 return False
1453
1454
1455 def _verify_report_set(
1456 report: schema.Report,
1457 entity_reports: list[tuple[str, schema.Report]] | None,
1458 *,
1459 allow_network: bool,
1460 ) -> None:
1461 reports = [item for _, item in entity_reports] if entity_reports else [report]
1462 for current_report in reports:
1463 freshness.verify_report(current_report, allow_network=allow_network)
1464 if not any(current_report.freshness_verdicts for current_report in reports):
1465 # An empty verdict list is a legitimate outcome, but a silent one has
1466 # already misled operators once; say why there is nothing to show.
1467 sys.stderr.write(
1468 "[last30days] Freshness verification found no re-checkable claims"
1469 " in this report; the verdict list is empty.\n"
1470 )
1471
1472
1473 def _run_cached_freshness(
1474 args: argparse.Namespace,
1475 config: dict[str, object],
1476 ) -> int:
1477 cached = _load_last_report_cache(
1478 None,
1479 ttl_seconds=_report_cache_ttl_seconds(config),
1480 )
1481 if cached is None:
1482 sys.stderr.write("[last30days] No fresh cached report; run a research pass first.\n")
1483 return 2
1484 report, entity_reports, cache_path = cached
1485 _verify_report_set(report, entity_reports, allow_network=not args.mock)
1486 if _update_cached_freshness(cache_path, report, entity_reports):
1487 sys.stderr.write(f"[last30days] Updated freshness verdicts in {cache_path}\n")
1488 else:
1489 sys.stderr.write("[last30days] warning: freshness cache update failed\n")
1490 return _render_save_and_print(args, report, entity_reports, None, config)
1491
1492
1493 def _drill_config(config: dict[str, object], sources: list[str]) -> dict[str, object]:
1494 """Enable configured comment enrichments for a deep follow-up."""
1495 drill_config = dict(config)
1496 include = {
1497 value.strip().lower()
1498 for value in str(config.get("INCLUDE_SOURCES") or "").split(",")
1499 if value.strip()
1500 }
1501 comment_flags = {
1502 "youtube": "youtube_comments",
1503 "tiktok": "tiktok_comments",
1504 "instagram": "instagram_comments",
1505 }
1506 include.update(comment_flags[source] for source in sources if source in comment_flags)
1507 if include:
1508 drill_config["INCLUDE_SOURCES"] = ",".join(sorted(include))
1509 drill_config["_drill_mode"] = True
1510 return drill_config
1511
1512
1513 def _run_drill(
1514 args: argparse.Namespace,
1515 config: dict[str, object],
1516 ) -> int:
1517 from lib import planner
1518
1519 cached = _load_last_report_cache(
1520 None,
1521 ttl_seconds=_report_cache_ttl_seconds(config),
1522 )
1523 if cached is None:
1524 sys.stderr.write(
1525 "[last30days] No fresh cached report; run a research pass first.\n"
1526 )
1527 return 2
1528 report, entity_reports, cache_path = cached
1529 if entity_reports:
1530 sys.stderr.write(
1531 "[last30days] Drill mode needs a single-topic cached report; "
1532 "run a research pass for one entity first.\n"
1533 )
1534 return 2
1535
1536 lookback_days = args.lookback_days
1537 if lookback_days is None:
1538 range_from = datetime.date.fromisoformat(report.range_from)
1539 range_to = datetime.date.fromisoformat(report.range_to)
1540 lookback_days = (range_to - range_from).days
1541 as_of_date = args.as_of_date or report.range_to
1542
1543 try:
1544 matched_clusters = planner.resolve_drill_clusters(report, args.drill)
1545 drill_plan = planner.build_drill_plan(
1546 report,
1547 args.drill,
1548 clusters=matched_clusters,
1549 )
1550 except planner.DrillTargetError as exc:
1551 sys.stderr.write(f"[last30days] {exc}\n")
1552 return 2
1553
1554 sources = list(drill_plan.source_weights)
1555 drill_config = _drill_config(config, sources)
1556 diag = pipeline.diagnose(drill_config, sources, safe=False)
1557 progress = ui.ProgressDisplay(
1558 f"{report.topic} — drill: {args.drill}",
1559 show_banner=True,
1560 )
1561 progress.start_processing()
1562 resolved = report.artifacts.get("resolved") or {}
1563 try:
1564 drill_report = pipeline.run(
1565 # Keep source gating anchored to the cached entity (for example,
1566 # StockTwits needs the original cashtag/finance context). The
1567 # external drill plan below remains cluster-focused.
1568 topic=report.topic,
1569 config=drill_config,
1570 depth="deep",
1571 requested_sources=sources,
1572 mock=args.mock,
1573 x_handle=(
1574 (args.x_handle or resolved.get("x_handle") or None)
1575 if "x" in sources else None
1576 ),
1577 x_related=(
1578 [value.strip() for value in args.x_related.split(",") if value.strip()]
1579 if (args.x_related and "x" in sources) else None
1580 ),
1581 web_backend=args.web_backend,
1582 external_plan=schema.to_dict(drill_plan),
1583 subreddits=(
1584 ([value.strip().removeprefix("r/") for value in args.subreddits.split(",") if value.strip()]
1585 if args.subreddits else list(resolved.get("subreddits") or []) or None)
1586 if "reddit" in sources else None
1587 ),
1588 tiktok_hashtags=(
1589 [value.strip().lstrip("#") for value in args.tiktok_hashtags.split(",") if value.strip()]
1590 if args.tiktok_hashtags else None
1591 ),
1592 tiktok_creators=(
1593 [value.strip().lstrip("@") for value in args.tiktok_creators.split(",") if value.strip()]
1594 if args.tiktok_creators else None
1595 ),
1596 ig_creators=(
1597 [value.strip().lstrip("@") for value in args.ig_creators.split(",") if value.strip()]
1598 if args.ig_creators else None
1599 ),
1600 lookback_days=lookback_days,
1601 as_of_date=as_of_date,
1602 github_user=(
1603 (args.github_user or resolved.get("github_user") or None)
1604 if "github" in sources else None
1605 ),
1606 github_repos=(
1607 ([value.strip() for value in args.github_repo.split(",") if value.strip()]
1608 if args.github_repo else list(resolved.get("github_repos") or []) or None)
1609 if "github" in sources else None
1610 ),
1611 trustpilot_domain=(
1612 (args.trustpilot_domain or resolved.get("trustpilot_domain") or None)
1613 if "trustpilot" in sources else None
1614 ),
1615 internal_subrun=True,
1616 corpus_dirs=args.corpus,
1617 corpus_all_time=args.corpus_all_time,
1618 )
1619 except Exception:
1620 progress.end_processing()
1621 raise
1622
1623 _show_runtime_ui(drill_report, progress, diag, suppress_web_promo=True)
1624 merged = pipeline.merge_drill_report(
1625 report,
1626 drill_report,
1627 matched_clusters,
1628 target=args.drill,
1629 )
1630 if _freshness_enabled(args, config):
1631 _verify_report_set(merged, None, allow_network=not args.mock)
1632 else:
1633 merged.freshness_verdicts = []
1634 if _write_last_run(report.topic, merged):
1635 sys.stderr.write(f"[last30days] Updated drill cache in {cache_path}\n")
1636 else:
1637 sys.stderr.write(
1638 "[last30days] warning: drill cache update failed; the next drill "
1639 "will see the pre-drill report\n"
1640 )
1641
1642 store_default = str(
1643 os.environ.get("LAST30DAYS_STORE")
1644 or config.get("LAST30DAYS_STORE")
1645 or ""
1646 ).lower()
1647 if args.store or store_default in {"1", "true", "yes"}:
1648 counts = persist_report(merged, store_db=_scoped_store_db(args))
1649 sys.stderr.write(
1650 f"[last30days] Stored {counts['new']} new, "
1651 f"{counts['updated']} updated findings\n"
1652 )
1653
1654 synthesis_md = None
1655 if args.synthesis_file:
1656 if args.emit == "html":
1657 synthesis_md = read_synthesis_file(args.synthesis_file)
1658 else:
1659 sys.stderr.write(
1660 "[last30days] Warning: --synthesis-file is only used with "
1661 "--emit=html; ignoring.\n"
1662 )
1663 return _render_save_and_print(args, merged, None, synthesis_md, config)
1664
1665
1666 def _save_discovery_output(
1667 rendered: str,
1668 *,
1669 domain: str,
1670 emit: str,
1671 save_dir: str,
1672 suffix: str = "",
1673 ) -> Path:
1674 directory = Path(save_dir).expanduser().resolve()
1675 directory.mkdir(parents=True, exist_ok=True)
1676 extension = "json" if emit == "json" else "md"
1677 safe_suffix = sanitize_suffix(suffix)
1678 suffix_part = f"-{safe_suffix}" if safe_suffix else ""
1679 stem = f"{slugify(domain)}-discover-raw{suffix_part}"
1680 date_str = datetime.datetime.now().strftime("%Y-%m-%d")
1681 candidates = [directory / f"{stem}.{extension}", directory / f"{stem}-{date_str}.{extension}"]
1682 candidates.extend(directory / f"{stem}-{date_str}-{index}.{extension}" for index in range(1, 100))
1683 encoded = rendered.encode("utf-8")
1684 for candidate in candidates:
1685 try:
1686 fd = os.open(candidate, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o644)
1687 except FileExistsError:
1688 continue
1689 with os.fdopen(fd, "wb") as output:
1690 output.write(encoded)
1691 return candidate
1692 raise RuntimeError("Could not find a unique discovery output filename")
1693
1694
1695 def _pre_run_prior_state(
1696 prior: dict[str, object] | None, run_ref: str
1697 ) -> dict[str, object] | None:
1698 """Reconstruct the queue state a topic had BEFORE this run identity
1699 recorded it.
1700
1701 A row whose last_run_ref equals THIS run's run_ref was stamped by this
1702 very run's own earlier attempt (a finalize retry), so its surface_count
1703 already includes this run's surfacing: subtract it and keep the prior's
1704 covered state (covered_at intact) so the retry renders exactly like the
1705 first attempt did. Only when nothing remains after the subtraction AND
1706 the row was never covered is the topic genuinely first-ever (no prior).
1707 """
1708 if not prior or prior.get("last_run_ref") != run_ref:
1709 return prior
1710 previously = max(0, int(prior["surface_count"]) - 1)
1711 if previously == 0 and prior["status"] != "covered":
1712 return None
1713 adjusted = dict(prior)
1714 adjusted["surface_count"] = previously
1715 return adjusted
1716
1717
1718 def _annotate_and_record_discovery_queue(
1719 report: schema.DiscoveryReport,
1720 args: argparse.Namespace,
1721 config: dict[str, object],
1722 run_ref: str | None = None,
1723 ) -> schema.DiscoveryReport:
1724 """Stamp queue annotations onto report topics, then record this surfacing.
1725
1726 Order matters: annotations describe the queue state BEFORE this run, so
1727 each topic is matched first and recorded second. The queue is on by
1728 default; the resolved config value LAST30DAYS_DISCOVERY_QUEUE == "off"
1729 (env var or .env, via env.get_config) disables it. Scoped runs
1730 (--save-dir) write the scoped research.db, never the global one. Runs
1731 synchronously after the pipeline returns - this writes disk, so the
1732 abandon-on-timeout daemon-thread pattern is forbidden here.
1733
1734 ``run_ref`` overrides the run identity: the finalize leg passes the
1735 pending report's leg-2 run_ref through so a finalize retry records (and
1736 annotates) as the SAME run - store.record_discovery_surfacing skips the
1737 double-count, and rows this very run identity stamped are not "prior"
1738 state, so retries render identically instead of claiming a resurfacing.
1739 """
1740 queue_setting = str(config.get("LAST30DAYS_DISCOVERY_QUEUE") or "").strip().lower()
1741 if queue_setting == "off" or not report.topics:
1742 return report
1743
1744 import dataclasses
1745
1746 import store
1747
1748 run_ref = run_ref or f"discover:{report.domain or 'trending'}:{report.generated_at}"
1749 as_of = (report.generated_at or "")[:10] or report.range_to
1750 annotated: list[schema.DiscoveryTopic] = []
1751 with store.scoped_db(_scoped_store_db(args)):
1752 store.init_db()
1753 # Phase 1: match EVERY topic before recording ANY. Interleaving
1754 # match+record in one loop lets topic N fuzzy-match a same-anchor
1755 # sibling row this very run recorded seconds earlier, falsely
1756 # annotating a first-ever topic as "surfaced 2nd time".
1757 # A row stamped by THIS run identity is this run's own earlier
1758 # attempt (finalize retry), not prior state: reconstruct the pre-run
1759 # state (count minus this run's own surfacing, covered state kept)
1760 # so retries render identically for topics WITH history too.
1761 priors = [
1762 _pre_run_prior_state(prior, run_ref)
1763 for prior in (
1764 store.match_discovery_topic(topic.name) for topic in report.topics
1765 )
1766 ]
1767 # Phase 2: record this run's surfacings. A topic whose (possibly
1768 # fuzzy) prior row is covered inherits that covered state, so a
1769 # user's covered mark survives judge naming drift instead of
1770 # silently forking into a fresh uncovered row.
1771 for topic, prior in zip(report.topics, priors):
1772 inherit_covered_at = None
1773 if prior and prior["status"] == "covered":
1774 inherit_covered_at = prior["covered_at"] or prior["last_surfaced"]
1775 store.record_discovery_surfacing(
1776 topic.name,
1777 domain=report.domain,
1778 run_ref=run_ref,
1779 as_of=as_of,
1780 inherit_covered_at=inherit_covered_at,
1781 )
1782 for topic, prior in zip(report.topics, priors):
1783 if prior:
1784 topic = dataclasses.replace(
1785 topic,
1786 previously_surfaced_count=prior["surface_count"],
1787 last_surfaced=prior["last_surfaced"],
1788 covered=prior["status"] == "covered",
1789 )
1790 annotated.append(topic)
1791 return dataclasses.replace(report, topics=annotated)
1792
1793
1794 def _record_discovery_queue_safely(
1795 report: schema.DiscoveryReport,
1796 args: argparse.Namespace,
1797 config: dict[str, object],
1798 run_ref: str | None = None,
1799 ) -> schema.DiscoveryReport:
1800 """Annotate + record the discovery queue, degrading a broken research.db
1801 (locked, read-only dir, corrupt) to a stderr warning: a queue failure
1802 must never destroy a finished pipeline run or the protocol's final
1803 brief. Shared verbatim by the one-shot and finalize paths."""
1804 try:
1805 return _annotate_and_record_discovery_queue(
1806 report, args, config, run_ref=run_ref,
1807 )
1808 except (sqlite3.Error, OSError) as exc:
1809 sys.stderr.write(
1810 f"[last30days] Warning: discovery queue unavailable ({exc}); "
1811 "continuing without queue annotations.\n"
1812 )
1813 return report
1814
1815
1816 def _emit_and_save_discovery_report(
1817 report: schema.DiscoveryReport,
1818 args: argparse.Namespace,
1819 domain: str,
1820 ) -> None:
1821 """Render a discovery report per --emit, honor --output/--save-dir, and
1822 print it. Shared verbatim by the one-shot and finalize paths."""
1823 if args.emit == "json":
1824 payload = schema.to_dict(report) if args.json_profile == "raw" else schema.to_discovery_export(report)
1825 rendered = json.dumps(payload, indent=2, sort_keys=True)
1826 else:
1827 rendered = render.render_discovery(report)
1828
1829 if args.output:
1830 output_path = save_rendered_output(rendered, args.output)
1831 sys.stderr.write(f"[last30days] Saved output to {output_path}\n")
1832 if args.save_dir:
1833 save_path = _save_discovery_output(
1834 rendered,
1835 domain=domain or "trending",
1836 emit=args.emit,
1837 save_dir=args.save_dir,
1838 suffix=args.save_suffix or "",
1839 )
1840 sys.stderr.write(f"[last30days] Saved output to {save_path}\n")
1841 print(rendered)
1842
1843
1844 def _discovery_strict_exit_code(
1845 source_status: dict[str, schema.SourceOutcome],
1846 config: dict[str, object],
1847 ) -> int:
1848 """The ONE LAST30DAYS_STRICT_EXIT evaluation for every discovery
1849 invocation - the one-shot and all three protocol legs (issue #384's
1850 discovery counterpart). Rendering/output already happened by the time
1851 this runs; only the exit code shifts to 3 when strict exit is on and any
1852 source outcome is neither clean nor an expected skip."""
1853 strict = str(config.get("LAST30DAYS_STRICT_EXIT") or "").strip().lower()
1854 if strict not in {"1", "true", "yes", "on"}:
1855 return 0
1856 degraded = sorted(
1857 source for source, outcome in (source_status or {}).items()
1858 if outcome.state not in _STRICT_EXIT_OK_STATES
1859 )
1860 if not degraded:
1861 return 0
1862 sys.stderr.write(
1863 f"[last30days] strict-exit: degraded sources: {', '.join(degraded)}\n"
1864 )
1865 sys.stderr.flush()
1866 return 3
1867
1868
1869 def _require_discover_mock_parity(
1870 loaded_mock: bool,
1871 args_mock: bool,
1872 *,
1873 label: str,
1874 path: Path | None,
1875 ) -> None:
1876 """A protocol leg's --mock flag must match the loaded handoff file's
1877 stamped provenance: mock-born state finalized by a real run would fake a
1878 real brief from fixture data, and real state finalized by --mock would
1879 silently drop the round's queue write. Mismatch is a contract failure
1880 (exit 2 via HandoffContractError)."""
1881 if bool(loaded_mock) == bool(args_mock):
1882 return
1883 location = str(path) if path is not None else "(unknown path)"
1884 if loaded_mock:
1885 raise discovery_handoff.HandoffContractError(
1886 f"{label} {location} is mock-born (a --mock leg wrote it): "
1887 "mock-born state cannot be finalized by a real run. Re-run this "
1888 "leg with --mock, or start a fresh real `--discover "
1889 "--nominate-only` sweep."
1890 )
1891 raise discovery_handoff.HandoffContractError(
1892 f"{label} {location} was written by a real run: real state "
1893 "cannot be finalized by a --mock run. Drop --mock, or start a fresh "
1894 "`--discover --nominate-only --mock` sweep."
1895 )
1896
1897
1898 def _run_queue_list(args: argparse.Namespace, config: dict[str, object]) -> int:
1899 """List uncovered surfaced topics from the persistent discovery queue."""
1900 import store
1901
1902 db_path = _scoped_store_db(args)
1903 if not Path(db_path or store.DB_PATH).exists():
1904 print("Discovery queue is empty - no discovery run has recorded topics yet.")
1905 return 0
1906 with store.scoped_db(db_path):
1907 rows = store.list_discovery_queue(status="surfaced")
1908 if not rows:
1909 # An existing db with zero queue rows (e.g. created via --store)
1910 # means no discovery run has recorded anything - only claim
1911 # "every topic is covered" when covered rows actually exist.
1912 if store.list_discovery_queue():
1913 print("Discovery queue is empty - every surfaced topic is marked covered.")
1914 else:
1915 print("Discovery queue is empty - no discovery run has recorded topics yet.")
1916 return 0
1917
1918 headers = ("name", "domain", "surface_count", "last_surfaced", "status")
1919 table = [
1920 (
1921 str(row["name"]),
1922 str(row["domain"] or "-"),
1923 str(row["surface_count"]),
1924 str(row["last_surfaced"]),
1925 str(row["status"]),
1926 )
1927 for row in rows
1928 ]
1929 widths = [
1930 max(len(headers[column]), *(len(row[column]) for row in table))
1931 for column in range(len(headers))
1932 ]
1933 lines = [
1934 " ".join(header.ljust(widths[i]) for i, header in enumerate(headers)).rstrip(),
1935 " ".join("-" * widths[i] for i in range(len(headers))),
1936 ]
1937 lines.extend(
1938 " ".join(row[i].ljust(widths[i]) for i in range(len(headers))).rstrip()
1939 for row in table
1940 )
1941 print("\n".join(lines))
1942 return 0
1943
1944
1945 def _run_queue_cover(
1946 args: argparse.Namespace,
1947 config: dict[str, object],
1948 name: str,
1949 ) -> int:
1950 """Mark a queued discovery topic covered; unknown names error loudly."""
1951 import store
1952
1953 if not name:
1954 sys.stderr.write(
1955 "[last30days] queue cover requires a topic name: "
1956 'queue cover "<topic name>".\n'
1957 )
1958 return 2
1959 db_path = _scoped_store_db(args)
1960 if not Path(db_path or store.DB_PATH).exists():
1961 sys.stderr.write(
1962 f"[last30days] No queued topic named {name!r}: the discovery queue "
1963 "is empty (no discovery run has recorded topics yet).\n"
1964 )
1965 return 2
1966 with store.scoped_db(db_path):
1967 row = store.mark_discovery_covered(
1968 name, as_of=datetime.date.today().isoformat()
1969 )
1970 if row is None:
1971 sys.stderr.write(
1972 f"[last30days] No queued topic named {name!r}. Covering requires "
1973 "the exact topic name; run 'queue list' to see queued names.\n"
1974 )
1975 return 2
1976 print(f"Marked covered: {row['name']} (covered {row['covered_at']})")
1977 return 0
1978
1979
1980 def _resolve_discovery_source_boundary(
1981 args: argparse.Namespace, config: dict[str, object],
1982 ) -> tuple[list[str] | None, list[str] | None] | None:
1983 """Resolve the discovery sweep's source lists from the user's boundary.
1984
1985 Returns ``(listing_sources, enrichment_boundary)`` - the discovery-capable
1986 subset for the sweep, and the user's ORIGINAL boundary honored by the
1987 per-topic research passes (which reach beyond the listing feeds - e.g.
1988 Techmeme, arXiv, YouTube, Polymarket); both None mean every available
1989 source. Returns None (after writing the exit-2 error) when the configured
1990 boundary leaves nothing to sweep: silently widening to all feeds would
1991 query sources the user filtered out.
1992 """
1993 requested_sources = resolve_requested_sources(args.search, config)
1994 enrich_requested_sources = list(requested_sources) if requested_sources else None
1995 if requested_sources:
1996 discovery_sources = [
1997 source for source in requested_sources
1998 if source in pipeline.DISCOVERY_SOURCES
1999 ]
2000 if not discovery_sources:
2001 origin = "--search" if args.search is not None else "LAST30DAYS_DEFAULT_SEARCH"
2002 sys.stderr.write(
2003 f"[last30days] {origin} has no discovery-capable sources "
2004 f"(unsupported: {', '.join(requested_sources)}); discovery "
2005 f"sweeps use: {', '.join(pipeline.DISCOVERY_SOURCES)}. Pass "
2006 "--search with one of those (or clear the source filter) to "
2007 "run a sweep.\n"
2008 )
2009 return None
2010 requested_sources = discovery_sources
2011 return requested_sources, enrich_requested_sources
2012
2013
2014 def _discover_subreddits(args: argparse.Namespace) -> list[str] | None:
2015 return (
2016 [value.strip().removeprefix("r/") for value in args.subreddits.split(",") if value.strip()]
2017 if args.subreddits else None
2018 )
2019
2020
2021 def _discover_domain(args: argparse.Namespace) -> str:
2022 """The whitespace-normalized discovery domain; empty = global trending."""
2023 return " ".join(str(args.discover or "").split())
2024
2025
2026 def _run_discover(args: argparse.Namespace, config: dict[str, object]) -> int:
2027 domain = _discover_domain(args)
2028 # Empty domain = global trending: sweep every river feed's hot list with no
2029 # keyword gate. The confidence floor is what keeps junk out, not a keyword.
2030 # (--as-of and HTML rejection live in _main's shared --discover dispatch,
2031 # so every leg - one-shot or protocol - applies the same guards.)
2032 if args.synthesis_file:
2033 sys.stderr.write("[last30days] Warning: --synthesis-file is not used by discovery mode.\n")
2034
2035 boundary = _resolve_discovery_source_boundary(args, config)
2036 if boundary is None:
2037 return 2
2038 requested_sources, enrich_requested_sources = boundary
2039 subreddits = _discover_subreddits(args)
2040 depth = "deep" if args.deep else "quick" if args.quick else "default"
2041 try:
2042 report = pipeline.run_discover(
2043 domain=domain,
2044 config=config,
2045 depth=depth,
2046 requested_sources=requested_sources,
2047 mock=args.mock,
2048 subreddits=subreddits,
2049 lookback_days=args.lookback_days or 30,
2050 as_of_date=args.as_of_date,
2051 enrich=not args.discover_shallow,
2052 enrich_requested_sources=enrich_requested_sources,
2053 )
2054 except ValueError as exc:
2055 sys.stderr.write(f"[last30days] {exc}\n")
2056 return 2
2057
2058 # Persistent topic queue: annotate this report from prior surfacings, then
2059 # record this run's surfacings - BEFORE rendering/export so the Pipeline
2060 # line and the JSON queue fields see the annotations. Mock runs stay 100%
2061 # side-effect-free.
2062 if not args.mock:
2063 report = _record_discovery_queue_safely(report, args, config)
2064
2065 _emit_and_save_discovery_report(report, args, domain)
2066 return _discovery_strict_exit_code(report.source_status, config)
2067
2068
2069 def _discover_handoff_state_dir(args: argparse.Namespace) -> Path | None:
2070 """One resolver for every protocol leg's handoff files: the save dir when
2071 given (mirroring _scoped_store_db's scoping), else the config dir - the
2072 same base _last_report_cache_path uses. args.save_dir is read AFTER the
2073 LAST30DAYS_MEMORY_DIR fallback in _main resolved it."""
2074 return discovery_handoff.handoff_state_dir(
2075 getattr(args, "save_dir", None), env.CONFIG_DIR
2076 )
2077
2078
2079 def _run_discover_nominate(args: argparse.Namespace, config: dict[str, object]) -> int:
2080 """Protocol leg 1: sweep the listings, build the full judge pool, write
2081 the nominations bundle, and print the host-facing judging digest.
2082
2083 No stage-1 judge, enrichment, confidence floor, or queue writes happen on
2084 this leg - the host judges from the bundle and leg 2 (--judgments)
2085 resumes from it. A zero-nomination sweep short-circuits to the existing
2086 nothing-solid brief with NO bundle written: there is nothing to judge.
2087 Writing a fresh bundle starts a NEW protocol round, so any pending
2088 report left by a prior round is deleted alongside it.
2089 """
2090 domain = _discover_domain(args)
2091 boundary = _resolve_discovery_source_boundary(args, config)
2092 if boundary is None:
2093 return 2
2094 requested_sources, enrich_requested_sources = boundary
2095 lookback_days = args.lookback_days or 30
2096 try:
2097 result = pipeline.run_discover_nominate(
2098 domain=domain,
2099 config=config,
2100 depth="deep" if args.deep else "quick" if args.quick else "default",
2101 requested_sources=requested_sources,
2102 mock=args.mock,
2103 subreddits=_discover_subreddits(args),
2104 lookback_days=lookback_days,
2105 as_of_date=args.as_of_date,
2106 )
2107 except ValueError as exc:
2108 sys.stderr.write(f"[last30days] {exc}\n")
2109 return 2
2110
2111 if not result.pool:
2112 print(render.render_discovery(pipeline.nominate_nothing_solid_report(result)))
2113 return _discovery_strict_exit_code(result.source_status, config)
2114
2115 entries = [
2116 discovery_handoff.PoolEntry(
2117 nomination=nomination,
2118 cluster_id=cluster_id,
2119 # No provider runs on this leg, so the nomination's name and junk
2120 # flag ARE the topic_shape heuristics - stored on the row as
2121 # leg 2's fallback for anything the host leaves unjudged.
2122 heuristic_name=nomination.name,
2123 heuristic_junk=nomination.junk_shape,
2124 )
2125 for nomination, cluster_id in result.pool
2126 ]
2127 bundle = discovery_handoff.write_nominations_bundle(
2128 entries,
2129 domain=result.plan.domain,
2130 tier="shallow" if args.discover_shallow else "deep",
2131 from_date=result.from_date,
2132 to_date=result.to_date,
2133 lookback_days=lookback_days,
2134 enrichment_source_boundary=enrich_requested_sources,
2135 requested_sources=requested_sources,
2136 # The sweep's finalized per-source outcomes ride the bundle so legs
2137 # 2-3 report degraded coverage instead of silently reading clean; the
2138 # mock stamp keeps mock-born and real state from cross-finalizing.
2139 source_status=result.source_status,
2140 mock=args.mock,
2141 # Same resolution as _discover_handoff_state_dir: save dir when
2142 # given, else the config dir.
2143 save_dir=getattr(args, "save_dir", None),
2144 config_dir=env.CONFIG_DIR,
2145 )
2146 # A fresh bundle starts a NEW protocol round: a pending report left by a
2147 # prior round is cross-round state a bare --finalize could silently
2148 # consume - delete it (missing file is a no-op).
2149 state_dir = _discover_handoff_state_dir(args)
2150 if state_dir is not None:
2151 discovery_handoff.pending_report_path(state_dir).unlink(missing_ok=True)
2152 print(discovery_handoff.build_host_digest(bundle))
2153 print(
2154 "\nJudgments file schema (leg 2): "
2155 f'{{"bundle_id": "{bundle.bundle_id}", "judgments": '
2156 '[{"id": "n1", "name": "<short topic name>", "junk": false, '
2157 '"worthiness": 0-100}, ...]}. '
2158 "Then resume with: --discover --judgments <path>."
2159 )
2160 return _discovery_strict_exit_code(result.source_status, config)
2161
2162
2163 def _run_discover_resume(args: argparse.Namespace, config: dict[str, object]) -> int:
2164 """Protocol leg 2: resume from the nominations bundle, apply the host
2165 judgments file, run the deep per-topic research pass, and persist the
2166 ranked result as the pending report for leg 3 (--finalize).
2167
2168 Contract failures (missing/stale bundle, judgments not bound to it, a
2169 bundle whose mock provenance disagrees with this run's --mock flag, an
2170 unwritable pending-report path) raise HandoffContractError and map to
2171 exit 2 in _run_discover_protocol_leg. Zero floor survivors renders the
2172 nothing-solid brief right here (clearing any stale prior-round pending
2173 file): no pending file, no leg 3. No queue writes and no artifact saves
2174 happen on this leg - the topic queue and the rendered brief belong to
2175 leg 3.
2176 """
2177 save_dir = getattr(args, "save_dir", None)
2178 bundle = discovery_handoff.read_nominations_bundle(
2179 save_dir=save_dir, config_dir=env.CONFIG_DIR,
2180 )
2181 _require_discover_mock_parity(
2182 bundle.mock, args.mock,
2183 label="Nominations bundle", path=bundle.path,
2184 )
2185 judgments = discovery_handoff.read_judgments(
2186 args.judgments, bundle, save_dir=save_dir, config_dir=env.CONFIG_DIR,
2187 )
2188 result = pipeline.run_discover_resume(
2189 bundle, judgments, config=config, mock=args.mock,
2190 )
2191 report = result.report
2192
2193 if not report.topics:
2194 # Nothing cleared the floor: the honest brief ends the protocol here.
2195 # This round wrote no pending file, so a stale one from an earlier
2196 # round must not survive to feed a bare --finalize (missing file is
2197 # a no-op).
2198 state_dir = _discover_handoff_state_dir(args)
2199 if state_dir is not None:
2200 discovery_handoff.pending_report_path(state_dir).unlink(missing_ok=True)
2201 print(render.render_discovery(report))
2202 return _discovery_strict_exit_code(report.source_status, config)
2203
2204 state_dir = _discover_handoff_state_dir(args)
2205 if state_dir is None:
2206 # Unreachable in practice - reading the bundle above required one of
2207 # these locations - but kept as a loud contract error, not an assert.
2208 raise discovery_handoff.HandoffContractError(
2209 "No handoff location available to write the pending report: "
2210 "pass --save-dir or configure ~/.config/last30days/."
2211 )
2212 pending_path = discovery_handoff.pending_report_path(state_dir)
2213 payload = {
2214 "kind": schema.DISCOVERY_PENDING_KIND,
2215 "schema_version": schema.DISCOVERY_PENDING_SCHEMA_VERSION,
2216 "bundle_id": bundle.bundle_id,
2217 # Fresh TTL clock: leg 3 measures staleness from THIS resume run,
2218 # not from the leg-1 sweep.
2219 "generated_at": report.generated_at,
2220 # Same run_ref format the queue records (leg 3 replays it verbatim).
2221 "run_ref": f"discover:{report.domain or 'trending'}:{report.generated_at}",
2222 # Leg-2 provenance: leg 3 refuses to finalize across the mock/real
2223 # boundary in either direction.
2224 "mock": bool(args.mock),
2225 # Full schema round-trip (the _write_last_run precedent): leg 3
2226 # rebuilds the report from this dict instead of re-running anything.
2227 "report": schema.to_dict(report),
2228 "angle_inputs": result.angle_inputs,
2229 }
2230 # ONE post-loop write from the main thread; enrichment workers are daemon
2231 # threads and never touch disk.
2232 try:
2233 state_dir.mkdir(parents=True, exist_ok=True)
2234 pending_path.write_text(json.dumps(payload, indent=2), encoding="utf-8")
2235 except OSError as exc:
2236 # A locked/read-only/full disk is the protocol's clean exit-2 path,
2237 # never a traceback (same contract as the bundle write).
2238 raise discovery_handoff.HandoffContractError(
2239 f"Could not write pending discovery report {pending_path}: {exc}"
2240 ) from exc
2241
2242 print(
2243 f"Judged discovery resume: {len(report.topics)} topic"
2244 f"{'s' if len(report.topics) != 1 else ''} cleared the floor "
2245 f"(bundle_id {bundle.bundle_id})."
2246 )
2247 print(f"Pending report: {pending_path}")
2248 print("\nAngle inputs by nomination id:")
2249 print(json.dumps(result.angle_inputs, indent=2))
2250 print(
2251 "\nWrite the angles file (leg 3): "
2252 f'{{"bundle_id": "{bundle.bundle_id}", "angles": '
2253 '[{"id": "n1", "podcast": "<one-sentence hook>", '
2254 '"x_article": "<one-sentence hook>"}, ...]} - one row per topic id '
2255 "above.\n"
2256 "Then finalize with: --discover --finalize --angles <path>."
2257 )
2258 return _discovery_strict_exit_code(report.source_status, config)
2259
2260
2261 def _run_discover_finalize(args: argparse.Namespace, config: dict[str, object]) -> int:
2262 """Protocol leg 3: load the leg-2 pending report, apply host angles,
2263 render the final brief, save discovery artifacts, and record the topic
2264 queue. The cheap offline leg - no sweep, no enrichment, no providers,
2265 no network; everything renders from the pending report. (HTML/--as-of
2266 rejection lives in _main's shared --discover dispatch.)
2267
2268 Contract failures (missing/stale/mismatched pending report or angles)
2269 raise HandoffContractError and map to exit 2 in
2270 _run_discover_protocol_leg. The pending file is deliberately LEFT IN
2271 PLACE on success: a finalize retry with a corrected angles file must
2272 keep working within the TTL, and the queue records under the pending
2273 report's leg-2 run_ref, so retries never double-count a surfacing.
2274 Mock finalize renders identically but writes no queue rows.
2275 """
2276 import dataclasses
2277
2278 save_dir = getattr(args, "save_dir", None)
2279 pending = discovery_handoff.read_pending_report(
2280 save_dir=save_dir, config_dir=env.CONFIG_DIR,
2281 )
2282 _require_discover_mock_parity(
2283 pending.mock, args.mock,
2284 label="Pending discovery report", path=pending.path,
2285 )
2286 angles = discovery_handoff.read_angles(
2287 args.angles, pending, save_dir=save_dir, config_dir=env.CONFIG_DIR,
2288 )
2289 try:
2290 report = schema.discovery_report_from_dict(pending.report)
2291 except (KeyError, TypeError, ValueError) as exc:
2292 # The envelope validated but the report body is structurally
2293 # incomplete: a contract failure with the resume remedy, never a
2294 # traceback out of the finalize leg.
2295 raise discovery_handoff.HandoffContractError(
2296 f"Pending discovery report {pending.path} carries a malformed "
2297 f"report body ({type(exc).__name__}: {exc}). "
2298 f"{discovery_handoff._RESUME_REMEDY}"
2299 ) from exc
2300
2301 if angles:
2302 # Host angles are keyed by nomination id; the pending report's
2303 # angle_inputs mapping carries each surviving id's applied topic
2304 # name, which is how angles land on the right DiscoveryTopic.
2305 angles_by_name = {
2306 name: host
2307 for nomination_id, host in angles.items()
2308 if (name := (pending.angle_inputs.get(nomination_id) or {}).get("name"))
2309 }
2310 report = dataclasses.replace(report, topics=[
2311 dataclasses.replace(
2312 topic,
2313 podcast_angle=host.podcast,
2314 x_article_angle=host.x_article,
2315 )
2316 if (host := angles_by_name.get(topic.name)) is not None
2317 else topic
2318 for topic in report.topics
2319 ])
2320
2321 # Persistent topic queue: the protocol's ONE queue write happens here,
2322 # under the leg-2 run identity (pending.run_ref) so finalize retries are
2323 # idempotent. Mock runs stay 100% side-effect-free.
2324 if not args.mock:
2325 report = _record_discovery_queue_safely(
2326 report, args, config, run_ref=pending.run_ref or None,
2327 )
2328
2329 _emit_and_save_discovery_report(report, args, report.domain)
2330 return _discovery_strict_exit_code(report.source_status, config)
2331
2332
2333 def _run_discover_protocol_leg(
2334 args: argparse.Namespace, config: dict[str, object]
2335 ) -> int:
2336 """Route one validated protocol invocation to its leg. Contract failures
2337 (unreadable/stale/mismatched handoff files) map to stderr + exit 2 here,
2338 so the leg bodies (U3-U5) raise HandoffContractError freely."""
2339 try:
2340 if args.nominate_only:
2341 return _run_discover_nominate(args, config)
2342 # --judgments dispatch keys on flag presence (is not None), matching
2343 # the --discover convention: never on the path string's truthiness.
2344 if args.judgments is not None:
2345 return _run_discover_resume(args, config)
2346 return _run_discover_finalize(args, config)
2347 except discovery_handoff.HandoffContractError as exc:
2348 sys.stderr.write(f"[last30days] {exc.message}\n")
2349 return 2
2350
2351
2352 _STRICT_EXIT_OK_STATES = {"ok", "no-results", "skipped-unconfigured"}
2353
2354
2355 def _strict_exit_code(
2356 report: schema.Report,
2357 entity_reports: list[tuple[str, schema.Report]] | None,
2358 config: dict[str, object],
2359 ) -> int:
2360 """Opt-in machine-detectable degraded-run signal (issue #384).
2361
2362 When LAST30DAYS_STRICT_EXIT is truthy, a run whose report carries any
2363 source outcome that is neither clean nor a plain no-results exits 3 so
2364 cron/CI wrappers can distinguish degraded coverage from success. Default
2365 behavior (exit 0, warning rendered in the report footer) is unchanged.
2366 """
2367 raw = str(config.get("LAST30DAYS_STRICT_EXIT") or "").strip().lower()
2368 if raw not in {"1", "true", "yes", "on"}:
2369 return 0
2370 reports = [report] + [rep for _, rep in (entity_reports or [])]
2371 degraded = sorted({
2372 name
2373 for rep in reports
2374 for name, outcome in (rep.source_status or {}).items()
2375 if outcome.state not in _STRICT_EXIT_OK_STATES
2376 })
2377 if not degraded:
2378 return 0
2379 sys.stderr.write(
2380 f"[last30days] strict-exit: degraded sources: {', '.join(degraded)}\n"
2381 )
2382 sys.stderr.flush()
2383 return 3
2384
2385
2386 def _audience_register_for_run(
2387 args: argparse.Namespace,
2388 config: dict[str, object],
2389 entity_reports: list[tuple[str, schema.Report]] | None,
2390 ) -> registers.AudienceRegister:
2391 """Resolve CLI > config for single-topic standard brief renderers."""
2392
2393 from lib import planner
2394
2395 topic = " ".join(getattr(args, "topic", [])).strip()
2396 comparison_topic_requested = bool(
2397 len(planner._comparison_entities(topic)) >= 2
2398 or args.competitors is not None
2399 or args.competitors_list
2400 or args.competitors_plan
2401 )
2402 if (
2403 entity_reports
2404 or comparison_topic_requested
2405 or args.drill
2406 or args.emit not in {"compact", "md", "html"}
2407 ):
2408 return registers.get_register()
2409 explicit = getattr(args, "register", None)
2410 configured = config.get("LAST30DAYS_REGISTER")
2411 name = explicit or (str(configured) if configured else "default")
2412 # Preserve configs written by the pre-register ELI5 follow-up command.
2413 legacy_eli5 = str(config.get("ELI5_MODE") or "").strip().lower()
2414 if not explicit and not configured and legacy_eli5 in {"1", "true", "yes", "on"}:
2415 name = "eli5"
2416 return registers.get_register(name)
2417
2418
2419 def _render_save_and_print(
2420 args: argparse.Namespace,
2421 report: schema.Report,
2422 entity_reports: list[tuple[str, schema.Report]] | None,
2423 synthesis_md: str | None,
2424 config: dict[str, object],
2425 ) -> int:
2426 fun_level = str(config.get("FUN_LEVEL", "medium")).lower()
2427 try:
2428 audience = _audience_register_for_run(args, config, entity_reports)
2429 except ValueError as exc:
2430 sys.stderr.write(f"[last30days] {exc}\n")
2431 return 2
2432 if audience.name != "default":
2433 sys.stderr.write(f"[last30days] Audience register: {audience.name}\n")
2434 sys.stderr.flush()
2435 # Comparison HTML is the one case where the saved file's title and content
2436 # have to be overridden away from the leading entity's report. Compute the
2437 # gate once so the footer-display and save-output paths can't disagree.
2438 is_comparison_html = bool(entity_reports) and args.emit == "html"
2439 footer_save_path = None
2440 if args.output:
2441 footer_save_path = compute_output_path_display(args.output)
2442 elif args.save_dir:
2443 save_topic_for_display = comparison_topic(entity_reports) if is_comparison_html else report.topic
2444 footer_save_path = compute_save_path_display(
2445 args.save_dir, save_topic_for_display, args.save_suffix or "", args.emit
2446 )
2447
2448 if entity_reports:
2449 rendered = emit_comparison_output(
2450 entity_reports,
2451 args.emit,
2452 fun_level=fun_level,
2453 save_path=footer_save_path,
2454 synthesis_md=synthesis_md,
2455 json_profile=args.json_profile,
2456 )
2457 else:
2458 rendered = emit_output(
2459 report,
2460 args.emit,
2461 fun_level=fun_level,
2462 save_path=footer_save_path,
2463 synthesis_md=synthesis_md,
2464 json_profile=args.json_profile,
2465 register=audience.name,
2466 )
2467 has_private_corpus = _report_has_private_corpus(report) or bool(
2468 entity_reports
2469 and any(_report_has_private_corpus(entity) for _label, entity in entity_reports)
2470 )
2471 private_saved_format = has_private_corpus
2472 publish_companion_paths: list[Path] = []
2473 if args.output:
2474 output_path = save_rendered_output(
2475 rendered,
2476 args.output,
2477 private=private_saved_format,
2478 )
2479 if args.emit == "html":
2480 publish_companion_paths.append(output_path)
2481 sys.stderr.write(f"[last30days] Saved output to {output_path}\n")
2482 sys.stderr.flush()
2483 if args.save_dir:
2484 # Save the main topic's raw file (single-entity or comparison main).
2485 # Bind the render to the path save_output actually allocates so the
2486 # saved report and stdout agree even when collision fallback is used.
2487 def _render_with_actual_path(actual_path: Path) -> str:
2488 nonlocal rendered
2489 display = compute_output_path_display(str(actual_path))
2490 if entity_reports:
2491 rendered = emit_comparison_output(
2492 entity_reports,
2493 args.emit,
2494 fun_level=fun_level,
2495 save_path=display,
2496 synthesis_md=synthesis_md,
2497 json_profile=args.json_profile,
2498 )
2499 else:
2500 rendered = emit_output(
2501 report,
2502 args.emit,
2503 fun_level=fun_level,
2504 save_path=display,
2505 synthesis_md=synthesis_md,
2506 json_profile=args.json_profile,
2507 register=audience.name,
2508 )
2509 if args.emit not in {"json", "html"} and not entity_reports:
2510 # Markdown saves keep the complete debug artifact (all clusters
2511 # and per-source items), matching the render_fn-less path in
2512 # save_output and the comparison peer saves. Saving the compact
2513 # stdout render instead made most collected evidence
2514 # unrecoverable from the raw file (#923). The stdout re-render
2515 # above still runs so the visible footer cites the real path,
2516 # and the saved artifact carries the same citation.
2517 return render.render_full(report, save_path=display)
2518 return rendered
2519
2520 save_path = save_output(
2521 report,
2522 args.emit,
2523 args.save_dir,
2524 suffix=args.save_suffix or "",
2525 synthesis_md=synthesis_md,
2526 topic_override=comparison_topic(entity_reports) if is_comparison_html else None,
2527 json_profile=args.json_profile,
2528 register=audience.name,
2529 private=private_saved_format,
2530 render_fn=_render_with_actual_path,
2531 )
2532 if args.emit == "html":
2533 publish_companion_paths.append(save_path)
2534 sys.stderr.write(f"[last30days] Saved output to {save_path}\n")
2535 comparison_peer_paths: list[Path] = []
2536 # Competitor / vs-mode: also save a per-entity raw file for each peer.
2537 # Matches historical vs-mode behavior (N passes -> N save files).
2538 if entity_reports and len(entity_reports) > 1:
2539 for label, entity_report in entity_reports[1:]:
2540 peer_path = save_output(
2541 entity_report, args.emit, args.save_dir,
2542 suffix=args.save_suffix or "",
2543 synthesis_md=synthesis_md,
2544 json_profile=args.json_profile,
2545 private=_report_has_private_corpus(entity_report),
2546 )
2547 comparison_peer_paths.append(peer_path)
2548 sys.stderr.write(f"[last30days] Saved output to {peer_path}\n")
2549 peers_display = ", ".join(str(path) for path in comparison_peer_paths)
2550 sys.stderr.write(
2551 f"[last30days] Comparison artifact set: main={save_path}; "
2552 f"peers={peers_display}\n"
2553 )
2554 sys.stderr.flush()
2555 if args.publish_html:
2556 try:
2557 has_private_corpus = "corpus" in report.source_status or bool(
2558 entity_reports
2559 and any("corpus" in entity.source_status for _label, entity in entity_reports)
2560 )
2561 publish_rendered = rendered
2562 if has_private_corpus:
2563 sys.stderr.write(
2564 "[last30days] Excluding local corpus evidence and synthesis from published HTML.\n"
2565 )
2566 if entity_reports:
2567 publish_rendered = emit_comparison_output(
2568 [
2569 (label, schema.without_sources(entity, {"corpus"}))
2570 for label, entity in entity_reports
2571 ],
2572 "html",
2573 fun_level=fun_level,
2574 save_path=footer_save_path,
2575 synthesis_md=None,
2576 json_profile=args.json_profile,
2577 )
2578 else:
2579 publish_rendered = emit_output(
2580 schema.without_sources(report, {"corpus"}),
2581 "html",
2582 fun_level=fun_level,
2583 save_path=footer_save_path,
2584 synthesis_md=None,
2585 json_profile=args.json_profile,
2586 register=audience.name,
2587 )
2588 publish_result = publish_rendered_html(
2589 publish_rendered,
2590 password=_publish_password_for_args(args, config),
2591 companion_paths=publish_companion_paths,
2592 )
2593 sys.stderr.write(f"[last30days] Published HTML to {publish_result['url']}\n")
2594 for warning in publish_result.get("_metadata_errors") or []:
2595 sys.stderr.write(f"[last30days] Publish metadata warning: {warning}\n")
2596 if publish_result.get("update_key"):
2597 sys.stderr.write(
2598 "[last30days] ht-ml.app returned an update key; not writing it "
2599 "to stdout, HTML, or publish metadata.\n"
2600 )
2601 sys.stderr.flush()
2602 except Exception as exc:
2603 sys.stderr.write(f"[last30days] HTML publish failed: {exc}\n")
2604 sys.stderr.flush()
2605 print(rendered)
2606 return _strict_exit_code(report, entity_reports, config)
2607
2608
2609 def _propagate_config_to_environ(config: dict[str, object]) -> None:
2610 """Push relevant env keys to os.environ so provider modules can read them.
2611
2612 The env.get_config() function reads from a .env file, but providers.py
2613 reads from os.environ directly. Without this, OPENAI_BASE_URL and
2614 XAI_BASE_URL overrides are silently ignored. This is a no-op for
2615 keys that are already set in process env.
2616 """
2617 for key in ("OPENAI_BASE_URL", "XAI_BASE_URL", "OPENROUTER_BASE_URL"):
2618 val = config.get(key)
2619 if val and not os.environ.get(key):
2620 os.environ[key] = val
2621
2622
2623 def _setup_allows_browser_cookies(args: argparse.Namespace, extra_argv: list[str]) -> bool:
2624 return (
2625 not args.no_browser_cookies
2626 and not args.diagnose
2627 and not args.preflight
2628 and "--allow-browser-cookies" in extra_argv
2629 )
2630
2631
2632 SETUP_PASSTHROUGH_FLAGS = {
2633 "--allow-browser-cookies",
2634 "--device-auth",
2635 "--github",
2636 "--github-start",
2637 "--github-poll",
2638 "--openclaw",
2639 "--store-key",
2640 }
2641
2642 STORE_KEY_FLAG = "--store-key"
2643
2644
2645 def _split_store_key(extra_argv: list[str]) -> tuple[bool, str, list[str]]:
2646 """Pull ``--store-key <NAME>`` / ``--store-key=<NAME>`` out of ``extra_argv``.
2647
2648 Returns ``(present, name, remaining)``. ``name`` is "" when the flag has
2649 no value; ``remaining`` is every other passthrough token, so the regular
2650 allowlist check still applies to them.
2651 """
2652 present = False
2653 name = ""
2654 remaining: list[str] = []
2655 i = 0
2656 while i < len(extra_argv):
2657 arg = extra_argv[i]
2658 if arg == STORE_KEY_FLAG:
2659 present = True
2660 if i + 1 < len(extra_argv) and not extra_argv[i + 1].startswith("-"):
2661 name = extra_argv[i + 1]
2662 i += 2
2663 continue
2664 i += 1
2665 continue
2666 if arg.startswith(STORE_KEY_FLAG + "="):
2667 present = True
2668 name = arg[len(STORE_KEY_FLAG) + 1:]
2669 i += 1
2670 continue
2671 remaining.append(arg)
2672 i += 1
2673 return present, name, remaining
2674
2675
2676 # One credential line: longer than any real token, short enough that a
2677 # misdirected stream on stdin cannot grow memory.
2678 STORE_KEY_MAX_BYTES = 64 * 1024
2679
2680
2681 def _run_store_key(name: str) -> int:
2682 """``setup --store-key <NAME>``: persist one allowlisted credential from stdin.
2683
2684 Reads exactly one line from stdin (bounded to ``STORE_KEY_MAX_BYTES``),
2685 strips whitespace, and writes it to the global ``.env`` as a 0o600 secret
2686 through ``setup_wizard.write_api_key``. An existing line for the same
2687 name is replaced, so a rejected credential can be rotated by running the
2688 command again. The value never reaches stdout or stderr: stdout carries
2689 ``NAME=****`` plus a JSON line ``{"persisted": bool, "key": NAME}``. A
2690 name outside ``env.KEYCHAIN_KEYS`` or an empty value exits 2 without
2691 echoing anything.
2692 """
2693 from lib import setup_wizard
2694
2695 if name not in env.KEYCHAIN_KEYS:
2696 # Do not enumerate the allowlist here: on an official-only host a
2697 # failure hint must not name the legacy credential keys.
2698 sys.stderr.write(
2699 "[last30days] setup --store-key: unknown or missing key name "
2700 "(must be a credential name the engine loads from its .env; "
2701 "see CONFIGURATION.md).\n"
2702 )
2703 return 2
2704 value = sys.stdin.readline(STORE_KEY_MAX_BYTES).strip()
2705 if not value:
2706 sys.stderr.write(
2707 f"[last30days] setup --store-key {name}: empty value on stdin; "
2708 "pipe the credential as a single line.\n"
2709 )
2710 return 2
2711 persisted = bool(
2712 setup_wizard.write_api_key(env.CONFIG_FILE, value, key_name=name, replace=True)
2713 )
2714 print(f"{name}=****")
2715 print(json.dumps({"persisted": persisted, "key": name}))
2716 return 0 if persisted else 1
2717
2718 SKILL_ONLY_FLAGS = {
2719 "--agent",
2720 }
2721
2722 # Doctor passthrough: `doctor --json` / `doctor --cached` mirror the setup
2723 # passthrough pattern (neither is a global parser flag; they only mean
2724 # something to doctor). `--cached` serves the stored doctor-cache.json report
2725 # within its TTL and falls through to a live run otherwise.
2726 DOCTOR_PASSTHROUGH_FLAGS = {
2727 "--json",
2728 "--cached",
2729 "--postmortem",
2730 "--probe",
2731 }
2732
2733
2734 def _looks_inline_json(value: str) -> bool:
2735 """True when a --x-posts argument is JSON text rather than a path."""
2736 stripped = value.strip()
2737 return stripped.startswith(("{", "[")) or "\n" in value
2738
2739
2740 def _comparison_requested(args: argparse.Namespace, topic: str) -> bool:
2741 """Whether this invocation is a comparison run (vs-topic or competitor flags)."""
2742 from lib import planner as _planner
2743
2744 return any(
2745 value is not None
2746 for value in (args.competitors, args.competitors_list, args.competitors_plan)
2747 ) or len(_planner._comparison_entities(topic, uncapped=True)) >= 2
2748
2749
2750 def _read_x_envelope(
2751 path: str,
2752 topic: str,
2753 args: argparse.Namespace,
2754 *,
2755 x_handle: str | None,
2756 x_related: list[str] | None,
2757 ) -> x_envelope.Envelope:
2758 """Validate a host-fetched X envelope against this run's window and topic."""
2759 from_date, to_date = dates.get_date_range(
2760 args.lookback_days or 30, as_of_date=args.as_of_date
2761 )
2762 return x_envelope.read(
2763 path,
2764 (from_date, to_date),
2765 topic,
2766 handles=[x_handle] if x_handle else [],
2767 related=[h for h in (x_related or []) if h and h.strip()],
2768 )
2769
2770
2771 def _attach_entity_envelopes(comp_plan: dict[str, dict], args: argparse.Namespace) -> None:
2772 """Validate every per-entity ``x_posts`` path in a --competitors-plan.
2773
2774 Each envelope is checked against its own entity (topic) and that entry's
2775 ``x_handle``/``x_related`` handles, and stored on the entry as
2776 ``_x_envelope`` for the entity sub-run. Raises EnvelopeContractError.
2777 """
2778 for entry in comp_plan.values():
2779 raw = entry.get("x_posts")
2780 if not raw:
2781 continue
2782 if not isinstance(raw, str) or _looks_inline_json(raw):
2783 raise x_envelope.EnvelopeContractError(
2784 f"--competitors-plan entry {entry.get('_name', '')!r}: x_posts must "
2785 "be a file path to a last30days-x-posts/1 envelope, never inline JSON. "
2786 "Rewrite the plan entry, or drop its x_posts field."
2787 )
2788 related = entry.get("x_related") if isinstance(entry.get("x_related"), list) else None
2789 entry["_x_envelope"] = _read_x_envelope(
2790 raw, str(entry.get("_name") or ""), args,
2791 x_handle=entry.get("x_handle") if isinstance(entry.get("x_handle"), str) else None,
2792 x_related=[str(h) for h in related] if related else None,
2793 )
2794
2795
2796 def _combine_envelope_digests(main_sha256: str | None, entity_sha256: dict[str, str]) -> str | None:
2797 """One digest binding the last-report cache to every envelope a run uses.
2798
2799 A single top-level envelope is bound by its own file digest; per-entity
2800 comparison envelopes are folded, name-sorted, into one digest. Both the
2801 cache write (validated envelopes) and the cache lookup (planned paths)
2802 must go through here so a comparison cache can be reused.
2803 """
2804 parts: list[str] = []
2805 if main_sha256:
2806 parts.append(main_sha256)
2807 for name in sorted(entity_sha256):
2808 parts.append(f"{name}:{entity_sha256[name]}")
2809 if not parts:
2810 return None
2811 if len(parts) == 1 and main_sha256:
2812 return main_sha256
2813 return hashlib.sha256("|".join(parts).encode("utf-8")).hexdigest()
2814
2815
2816 def _x_envelope_digest(
2817 main: x_envelope.Envelope | None, comp_plan: dict[str, dict] | None
2818 ) -> str | None:
2819 """Digest of the validated envelopes this run used (cache write side)."""
2820 entity_sha256 = {
2821 name: entry["_x_envelope"].sha256
2822 for name, entry in (comp_plan or {}).items()
2823 if entry.get("_x_envelope") is not None
2824 }
2825 return _combine_envelope_digests(main.sha256 if main is not None else None, entity_sha256)
2826
2827
2828 def _validate_extra_argv(parser: argparse.ArgumentParser, topic: str, extra_argv: list[str]) -> None:
2829 if not extra_argv:
2830 return
2831 if topic.lower() == "setup":
2832 # --store-key carries a value token; the name itself is allowlisted
2833 # later in _run_store_key, not here.
2834 _, _, extra_argv = _split_store_key(extra_argv)
2835 unsupported = [arg for arg in extra_argv if arg not in SETUP_PASSTHROUGH_FLAGS]
2836 if unsupported:
2837 parser.error(
2838 "unsupported setup argument(s): "
2839 + ", ".join(unsupported)
2840 + f"; supported setup passthrough flags are {', '.join(sorted(SETUP_PASSTHROUGH_FLAGS))}"
2841 )
2842 return
2843 if topic.lower() == "doctor":
2844 unsupported = [arg for arg in extra_argv if arg not in DOCTOR_PASSTHROUGH_FLAGS]
2845 if unsupported:
2846 parser.error(
2847 "unsupported doctor argument(s): "
2848 + ", ".join(unsupported)
2849 + f"; supported doctor passthrough flags are {', '.join(sorted(DOCTOR_PASSTHROUGH_FLAGS))}"
2850 )
2851 return
2852 skill_only = [arg for arg in extra_argv if arg in SKILL_ONLY_FLAGS]
2853 other_unknown = [arg for arg in extra_argv if arg not in SKILL_ONLY_FLAGS]
2854 if skill_only:
2855 message = (
2856 "unsupported Python CLI argument(s): "
2857 + ", ".join(skill_only)
2858 + "; these are skill arguments and must not be forwarded to scripts/last30days.py"
2859 )
2860 if other_unknown:
2861 message += "; also unsupported: " + ", ".join(other_unknown)
2862 parser.error(message)
2863 parser.error("unsupported Python CLI argument(s): " + ", ".join(extra_argv))
2864
2865
2866 def _config_policy_for_args(args: argparse.Namespace, topic: str, extra_argv: list[str]) -> env.ConfigLoadPolicy:
2867 normalized_topic = topic.lower()
2868 is_library_command = (
2869 normalized_topic == "library feed"
2870 or normalized_topic == "library search"
2871 or normalized_topic.startswith("library search ")
2872 )
2873 # Queue commands are local SQLite reads/writes: like library commands they
2874 # must never trigger browser-cookie extraction or Keychain prompts.
2875 is_queue_command = (
2876 normalized_topic == "queue list"
2877 or normalized_topic == "queue cover"
2878 or normalized_topic.startswith("queue cover ")
2879 )
2880 is_cached_verification = bool(getattr(args, "verify_freshness", None)) and not normalized_topic
2881 if args.no_browser_cookies:
2882 browser_mode = "off"
2883 elif (
2884 args.diagnose or args.preflight or normalized_topic == "doctor"
2885 or is_library_command or is_queue_command or is_cached_verification
2886 ):
2887 # doctor is plan-only like --diagnose: it must never read cookies.
2888 # Cache-only freshness verification hits only point APIs (Polymarket,
2889 # GitHub, StockTwits) - no cookie-backed source, so no Keychain prompt.
2890 browser_mode = "plan_only"
2891 elif normalized_topic == "setup":
2892 browser_mode = "read" if _setup_allows_browser_cookies(args, extra_argv) else "off"
2893 else:
2894 browser_mode = "read"
2895 return env.ConfigLoadPolicy(
2896 browser_cookies=browser_mode,
2897 inspect_ignored_project_config=args.diagnose or args.preflight or normalized_topic == "doctor",
2898 )
2899
2900
2901 def _run_library_feed(args: argparse.Namespace, config: dict[str, object]) -> int:
2902 """Generate the local research index/feed and optionally publish it."""
2903 from lib import feed, html_publish, library
2904
2905 if args.publish_html:
2906 sys.stderr.write(
2907 "[last30days] library feed uses --publish, not --publish-html.\n"
2908 )
2909 return 2
2910 if args.output:
2911 sys.stderr.write(
2912 "[last30days] library feed writes index.html and feed.xml to --save-dir; "
2913 "--output is not supported.\n"
2914 )
2915 return 2
2916
2917 memory_dir = Path(args.save_dir).expanduser() if args.save_dir else library.DEFAULT_MEMORY_DIR
2918 output_dir = memory_dir.resolve()
2919 # Scoped libraries (--save-dir) must not mix in the global briefing
2920 # archive: a client-specific or publishable feed pulling unrelated default
2921 # briefings could publish them publicly. The default library keeps the
2922 # archive; a scoped one reads only its own directory.
2923 briefs_dir = (
2924 library.DEFAULT_BRIEFS_DIR if not args.save_dir else memory_dir / "briefings"
2925 )
2926 entries, notes = library.scan_library(memory_dir, briefs_dir)
2927 feed_author = str(
2928 config.get("LAST30DAYS_LIBRARY_OWNER") or "last30days research library"
2929 )
2930 output_dir.mkdir(parents=True, exist_ok=True)
2931 library_id = library.get_or_create_library_id(output_dir)
2932 rendered_briefs_dir = output_dir / "briefs"
2933 has_private_entries = any(
2934 render.PRIVATE_CORPUS_START in entry.content for entry in entries
2935 )
2936 _ensure_output_directory(rendered_briefs_dir, private=has_private_entries)
2937
2938 def _preserve_hand_written_page(existing_path: Path, generated_marker: str) -> None:
2939 """Back up any page library feed did not generate before overwriting it."""
2940 if not existing_path.exists():
2941 return
2942 try:
2943 marker_found = generated_marker in existing_path.read_text(encoding="utf-8")
2944 except (OSError, UnicodeDecodeError):
2945 marker_found = False
2946 if marker_found:
2947 return
2948 backup = existing_path.with_suffix(existing_path.suffix + ".bak")
2949 counter = 1
2950 while backup.exists():
2951 backup = existing_path.with_suffix(f"{existing_path.suffix}.bak{counter}")
2952 counter += 1
2953 existing_path.replace(backup)
2954 sys.stderr.write(
2955 f"[last30days] {existing_path.name} was not generated by "
2956 f"library feed; preserved the original at {backup.name}\n"
2957 )
2958
2959 publishable_brief_documents: dict[str, str] = {}
2960 for entry in entries:
2961 rendered = html_render.render_library_brief(entry)
2962 target = rendered_briefs_dir / entry.output_name
2963 _preserve_hand_written_page(target, html_render.LIBRARY_BRIEF_MARKER)
2964 save_rendered_output(
2965 rendered,
2966 str(target),
2967 private=render.PRIVATE_CORPUS_START in entry.content,
2968 )
2969 publishable_brief_documents[entry.entry_id] = html_render.render_library_brief(
2970 entry, include_private=False
2971 )
2972
2973 current_brief_names = {entry.output_name for entry in entries}
2974 for path in rendered_briefs_dir.glob("*.html"):
2975 is_orphan = path.name not in current_brief_names
2976 if not (is_orphan and library.is_generated_brief_name(path.name)):
2977 continue
2978 # A generated-looking name is not proof of ownership; only prune
2979 # pages that carry the renderer's own marker.
2980 try:
2981 generated = html_render.LIBRARY_BRIEF_MARKER in path.read_text(
2982 encoding="utf-8"
2983 )
2984 except (OSError, UnicodeDecodeError):
2985 generated = False
2986 if generated:
2987 path.unlink()
2988
2989 feed_xml = feed.render_atom(entries, library_id=library_id, author=feed_author)
2990 index_html = html_render.render_library_index(entries)
2991 feed_path = output_dir / "feed.xml"
2992 index_path = output_dir / "index.html"
2993 _preserve_hand_written_page(feed_path, "urn:last30days:research-library")
2994 _preserve_hand_written_page(
2995 index_path, "Generated locally by <strong>last30days</strong>"
2996 )
2997 feed_path.write_text(feed_xml, encoding="utf-8")
2998 index_path.write_text(index_html, encoding="utf-8")
2999
3000 for note in notes:
3001 sys.stderr.write(f"[last30days] Library note: {note}\n")
3002 sys.stderr.write(
3003 f"[last30days] Library feed generated {len(entries)} brief(s): "
3004 f"{index_path} and {feed_path}\n"
3005 )
3006
3007 if args.publish:
3008 password = _publish_password_for_args(args, config)
3009 entry_urls: dict[str, str] = {}
3010 try:
3011 brief_results = html_publish.publish_html_documents(
3012 publishable_brief_documents,
3013 password=password,
3014 )
3015 entry_urls = {
3016 entry_id: str(result["url"])
3017 for entry_id, result in brief_results.items()
3018 }
3019 if batch_error := getattr(brief_results, "error", None):
3020 raise batch_error
3021 published_index = html_render.render_library_index(
3022 entries,
3023 entry_urls=entry_urls,
3024 feed_url=None,
3025 )
3026 index_result = html_publish.publish_html(published_index, password=password)
3027 index_url = str(index_result["url"])
3028 except (html_publish.HtmlPublishError, KeyError, OSError) as exc:
3029 sys.stderr.write(f"[last30days] Library publish failed: {exc}\n")
3030 if entry_urls:
3031 sys.stderr.write(
3032 f"[last30days] Partial publish: {len(entry_urls)} public brief "
3033 "page(s) were created before the failure.\n"
3034 )
3035 return 1
3036
3037 # Keep the local artifacts useful as a record of the live publication.
3038 feed_path.write_text(
3039 feed.render_atom(
3040 entries,
3041 library_id=library_id,
3042 entry_urls=entry_urls,
3043 author=feed_author,
3044 ),
3045 encoding="utf-8",
3046 )
3047 index_path.write_text(
3048 html_render.render_library_index(entries, entry_urls=entry_urls),
3049 encoding="utf-8",
3050 )
3051 sys.stderr.write(f"[last30days] Published library to {index_url}\n")
3052 sys.stderr.write(f"[last30days] Local Atom feed: {feed_path}\n")
3053 print(
3054 f"Library: {index_url}\nFeed: {feed_path}\n"
3055 "Atom feed is local; host feed.xml on any static host (for example, GitHub Pages) "
3056 "to make it subscribable."
3057 )
3058 return 0
3059
3060 print(
3061 f"Library: {index_path}\nFeed: {feed_path}\n"
3062 "Atom feed is local; host feed.xml on any static host (for example, GitHub Pages) "
3063 "to make it subscribable."
3064 )
3065 return 0
3066
3067
3068 def _run_library_search(
3069 args: argparse.Namespace,
3070 config: dict[str, object],
3071 query: str,
3072 ) -> int:
3073 """Search saved briefs and store sightings without network access."""
3074 from lib import library, library_index
3075
3076 if not query.strip():
3077 sys.stderr.write("[last30days] library search requires a non-empty query.\n")
3078 return 2
3079 if args.publish or args.publish_html:
3080 sys.stderr.write("[last30days] library search does not publish output.\n")
3081 return 2
3082 if args.emit != "compact":
3083 sys.stderr.write("[last30days] library search currently supports text output only.\n")
3084 return 2
3085 if args.output:
3086 sys.stderr.write(
3087 "[last30days] library search prints to stdout; --output is not supported.\n"
3088 )
3089 return 2
3090
3091 memory_dir = Path(args.save_dir).expanduser() if args.save_dir else library.DEFAULT_MEMORY_DIR
3092 try:
3093 matches, synced = library_index.sync_and_search(
3094 query,
3095 memory_dir=memory_dir,
3096 briefs_dir=(
3097 memory_dir / "briefings" if args.save_dir else library.DEFAULT_BRIEFS_DIR
3098 ),
3099 db_path=(
3100 memory_dir.resolve() / ".last30days-library.db"
3101 if args.save_dir else library_index.DEFAULT_LIBRARY_DB
3102 ),
3103 # A scoped search must never merge in the shared store: one
3104 # client's sightings would leak into another client's scope. A
3105 # scoped store is read only if it exists inside the save dir.
3106 store_db_path=(
3107 memory_dir.resolve() / "research.db"
3108 if args.save_dir else library_index.DEFAULT_STORE_DB
3109 ),
3110 )
3111 except library_index.LibrarySearchUnavailable as exc:
3112 sys.stderr.write(f"[last30days] Library search unavailable: {exc}.\n")
3113 return 2
3114 except (OSError, sqlite3.DatabaseError) as exc:
3115 sys.stderr.write(f"[last30days] Library search failed: {exc}.\n")
3116 return 1
3117 for note in synced.notes:
3118 sys.stderr.write(f"[last30days] Library note: {note}\n")
3119 if synced.rebuilt:
3120 sys.stderr.write("[last30days] Rebuilt a corrupt library search index.\n")
3121 print(render.render_library_search(query, matches), end="")
3122 return 0
3123
3124
3125 def _looks_like_entity_topic(topic: str) -> bool:
3126 """Whether a topic names a person, company, or product rather than a theme.
3127
3128 Keys on brevity, not capitalization. People type lowercase: "bentgo",
3129 "peter steinberger" and "getenergy.com" are entity searches every bit as
3130 much as their title-cased forms, and requiring a capital meant the most
3131 common real-world spelling never resolved a handle.
3132
3133 A short topic is an entity search; a longer one is a theme. "Peter
3134 Steinberger", "bentgo" and "getenergy.com" qualify; "best AI coding tools
3135 2026" and "how to build agents that scale" do not. Question-shaped topics
3136 are themes regardless of length.
3137
3138 Used only to decide whether resolving an X handle is worth one web search,
3139 so a false negative costs the old behavior and a false positive costs a
3140 single search.
3141 """
3142 text = (topic or "").strip()
3143 if not text or text.endswith("?"):
3144 return False
3145 words = [w for w in re.findall(r"[A-Za-z0-9_.@'-]+", text) if w]
3146 if not words or len(words) > 4:
3147 return False
3148 if any(w.startswith("@") for w in words):
3149 return True
3150 # A theme reads as a phrase built from common words; an entity does not.
3151 common = {
3152 "best", "top", "how", "why", "what", "when", "vs", "versus", "guide",
3153 "tips", "review", "reviews", "news", "latest", "update", "updates",
3154 "trends", "tools", "and", "or", "for", "the", "with", "about",
3155 }
3156 return not any(w.lower() in common for w in words)
3157
3158
3159 def main() -> int:
3160 parser = build_parser()
3161 # Use parse_known_args so setup sub-flags (--device-auth, --github,
3162 # --openclaw) pass through without argparse hard-exiting.
3163 args, extra_argv = parser.parse_known_args()
3164 if args.record_fixtures:
3165 with http.recording_requests(Path(args.record_fixtures)):
3166 return _main(parser, args, extra_argv)
3167 return _main(parser, args, extra_argv)
3168
3169
3170 def _main(
3171 parser: argparse.ArgumentParser,
3172 args: argparse.Namespace,
3173 extra_argv: list[str],
3174 ) -> int:
3175 if args.debug:
3176 os.environ["LAST30DAYS_DEBUG"] = "1"
3177
3178 if args.welcome:
3179 from lib import setup_wizard
3180 print(setup_wizard.render_welcome())
3181 return 0
3182
3183 topic = " ".join(args.topic).strip()
3184 original_topic = topic
3185 _validate_extra_argv(parser, topic, extra_argv)
3186 if args.resolve_save_dir:
3187 print(env.resolve_memory_dir(args.save_dir))
3188 return 0
3189 if args.x_posts is not None and _looks_inline_json(args.x_posts):
3190 sys.stderr.write(
3191 "[last30days] --x-posts accepts a file path only (inline JSON is not "
3192 "accepted); write the envelope to a .json file and pass its path.\n"
3193 )
3194 return 2
3195 if args.publish and topic.lower() != "library feed":
3196 sys.stderr.write(
3197 "[last30days] --publish is only supported by the 'library feed' command.\n"
3198 )
3199 return 2
3200 if topic.lower() == "setup":
3201 # Persisting a credential needs no config load (no Keychain / pass
3202 # probes, no cookie policy), so it dispatches before get_config.
3203 store_key_present, store_key_name, _ = _split_store_key(extra_argv)
3204 if store_key_present:
3205 return _run_store_key(store_key_name)
3206
3207 config = env.get_config(policy=_config_policy_for_args(args, topic, extra_argv))
3208 if args.perplexity_search_type is not None:
3209 config["LAST30DAYS_PERPLEXITY_SEARCH_TYPE"] = args.perplexity_search_type
3210 # One memo per command: comparison mode runs pipeline.run per entity in
3211 # parallel, so the reset must not live inside the pipeline.
3212 http.reset_reddit_keyless_memo()
3213 reddit.reset_scrapecreators_memo()
3214 resolved_corpus_dirs = corpus.resolve_directories(
3215 args.corpus, config.get("LAST30DAYS_CORPUS_DIRS")
3216 )
3217 # EXCLUDE_SOURCES=corpus disables corpus retrieval entirely; the hosted
3218 # privacy bypass below must use the same predicate, or hosted users with
3219 # configured-but-excluded dirs silently lose the remote backend.
3220 excluded_sources = {
3221 value.strip().lower()
3222 for value in str(config.get("EXCLUDE_SOURCES") or "").split(",")
3223 if value.strip()
3224 }
3225 if "corpus" in excluded_sources:
3226 resolved_corpus_dirs = []
3227 if resolved_corpus_dirs:
3228 config["_CORPUS_DIRS"] = [str(path) for path in resolved_corpus_dirs]
3229 if _config_truthy(config.get("LAST30DAYS_CORPUS_IN_EXPORT")):
3230 config["_CORPUS_IN_EXPORT"] = True
3231 _propagate_config_to_environ(config)
3232
3233 # Env-var fallback for --save-dir, mirroring the LAST30DAYS_STORE pattern below.
3234 # Uses `is None` / `is not None` checks (not truthy `or`) at every layer so that
3235 # `--save-dir ""`, `LAST30DAYS_MEMORY_DIR=""` (shell-export-empty), and explicit
3236 # absence each correctly suppress save. An `or` chain would collapse the empty
3237 # shell-export into the same path as unset, silently falling through to .env.
3238 if args.save_dir is None:
3239 env_val = os.environ.get("LAST30DAYS_MEMORY_DIR")
3240 args.save_dir = env_val if env_val is not None else config.get("LAST30DAYS_MEMORY_DIR")
3241
3242 # Surface SSH-routing config as an env var so library modules (e.g.
3243 # youtube_yt) can read it without taking a config dependency. This
3244 # routes yt-dlp through `ssh <host>` to bypass YouTube's bot-wall on
3245 # datacenter IPs (see lib/youtube_yt.py for details).
3246 if config.get("LAST30DAYS_YOUTUBE_SSH_HOST") and "LAST30DAYS_YOUTUBE_SSH_HOST" not in os.environ:
3247 os.environ["LAST30DAYS_YOUTUBE_SSH_HOST"] = config["LAST30DAYS_YOUTUBE_SSH_HOST"]
3248
3249 if args.preflight:
3250 requested_sources = resolve_requested_sources(args.search, config)
3251 diag = pipeline.diagnose(config, requested_sources, safe=True)
3252 if args.save_dir or args.preflight_report_on_save_dir:
3253 preflight = permission_preflight.build(
3254 config,
3255 diag,
3256 planned_save_dir=args.save_dir,
3257 report_on_save_dir=args.preflight_report_on_save_dir,
3258 )
3259 else:
3260 preflight = diag["permission_preflight"]
3261 if args.emit == "json":
3262 print(json.dumps(preflight, indent=2, sort_keys=True))
3263 else:
3264 print(permission_preflight.render_text(preflight), end="")
3265 return 0
3266
3267 # Handle doctor subcommand: topic-word dispatch mirroring setup (exact
3268 # match only, so multi-word research topics containing "doctor" still
3269 # research normally). Aggregates probes/descriptors/prescriptions into
3270 # one grouped health surface; always exits 0.
3271 if topic.lower() == "doctor":
3272 from lib import doctor
3273 return doctor.run(
3274 config,
3275 emit_json=(args.emit == "json" or "--json" in extra_argv),
3276 cached="--cached" in extra_argv,
3277 postmortem="--postmortem" in extra_argv,
3278 probe="--probe" in extra_argv,
3279 )
3280
3281 if topic.lower() == "library feed":
3282 return _run_library_feed(args, config)
3283 if topic.lower() == "library search" or topic.lower().startswith("library search "):
3284 return _run_library_search(args, config, topic[len("library search") :].strip())
3285
3286 if topic.lower() == "queue list":
3287 return _run_queue_list(args, config)
3288 if topic.lower() == "queue cover" or topic.lower().startswith("queue cover "):
3289 return _run_queue_cover(args, config, topic[len("queue cover") :].strip())
3290
3291 # Handle setup subcommand
3292 if topic.lower() == "setup":
3293 from lib import setup_wizard
3294 if "--openclaw" in extra_argv:
3295 results = setup_wizard.run_openclaw_setup(config)
3296 print(json.dumps(results))
3297 return 0
3298 if any(f in extra_argv for f in ("--github", "--device-auth", "--github-start", "--github-poll")):
3299 if "--github-start" in extra_argv:
3300 results = setup_wizard.run_github_start()
3301 elif "--github-poll" in extra_argv:
3302 results = setup_wizard.run_github_poll()
3303 elif "--github" in extra_argv:
3304 results = setup_wizard.run_github_auth()
3305 else:
3306 results = setup_wizard.run_full_device_auth()
3307 # Persist the returned key so the paid sources activate on the next
3308 # run, and mask it in stdout so the secret never lands in the host
3309 # model's captured Bash output.
3310 api_key = results.get("api_key")
3311 status = results.get("status")
3312 if api_key:
3313 if status == "success":
3314 results["persisted"] = setup_wizard.write_api_key(env.CONFIG_FILE, api_key)
3315 elif status == "already_registered":
3316 results["persisted"] = True # key was already saved
3317 else:
3318 results.setdefault("persisted", False)
3319 # Mask for EVERY status that carries a key, not just success, so
3320 # the raw secret never reaches the host model's captured stdout.
3321 results["api_key"] = setup_wizard.mask_api_key(api_key)
3322 else:
3323 results["persisted"] = False
3324 print(json.dumps(results))
3325 return 0
3326 sys.stderr.write("Running auto-setup...\n")
3327 results = setup_wizard.run_auto_setup(
3328 config,
3329 allow_browser_cookies=_setup_allows_browser_cookies(args, extra_argv),
3330 )
3331 # Keep only successful browsers, including distinct service winners;
3332 # "auto" would also probe browsers that did not supply any cookies.
3333 found_browsers = dict.fromkeys(
3334 "firefox" if browser == "firefox-wsl" else browser
3335 for browser in results.get("cookies_found", {}).values()
3336 )
3337 from_browser = ",".join(found_browsers) or None
3338 results["env_written"] = setup_wizard.write_setup_config(
3339 env.CONFIG_FILE,
3340 from_browser=from_browser,
3341 browser_consent=(
3342 None if args.diagnose else _setup_allows_browser_cookies(args, extra_argv)
3343 ),
3344 )
3345 if not results["env_written"]:
3346 sys.stderr.write("Setup configuration could not be fully saved; some settings may already be saved.\n")
3347 return 1
3348 sys.stderr.write(setup_wizard.get_setup_status_text(results) + "\n")
3349 return 0
3350
3351 # Bare --discover (no domain) is global trending, so the dispatch keys on
3352 # "flag present" (is not None), never on the domain string's truthiness.
3353 if args.deep_research and not topic:
3354 sys.stderr.write(
3355 "[last30days] --deep-research requires a normal positional topic; "
3356 "it cannot be combined with discovery, drill, or cached-only modes.\n"
3357 )
3358 return 2
3359
3360 if args.discover is not None:
3361 if topic:
3362 sys.stderr.write(
3363 "[last30days] --discover supplies the domain and cannot be combined "
3364 "with a positional topic.\n"
3365 )
3366 return 2
3367 if args.drill:
3368 sys.stderr.write("[last30days] --discover and --drill are mutually exclusive.\n")
3369 return 2
3370 # Shared guards for EVERY discover invocation - the one-shot and all
3371 # three protocol legs - hoisted here so no leg can drift: discovery
3372 # sweeps live listings (never --as-of) and has no HTML pipeline yet.
3373 if args.as_of_date:
3374 sys.stderr.write(
3375 "[last30days] --as-of cannot be used with --discover because discovery "
3376 "sweeps current live listings.\n"
3377 )
3378 return 2
3379 if args.emit == "html" or args.publish_html:
3380 sys.stderr.write("[last30days] discovery mode does not support HTML publishing yet.\n")
3381 return 2
3382 # The three protocol legs are one-leg-per-invocation: each pairing
3383 # below asks for two legs at once, so name the combination and stop.
3384 # (--judgments/--angles dispatch on presence, never path truthiness.)
3385 for first, second, conflict in (
3386 ("--nominate-only", "--judgments", args.nominate_only and args.judgments is not None),
3387 ("--nominate-only", "--finalize", args.nominate_only and args.finalize),
3388 ("--judgments", "--finalize", args.judgments is not None and args.finalize),
3389 ):
3390 if conflict:
3391 sys.stderr.write(
3392 f"[last30days] {first} and {second} are mutually exclusive: "
3393 "each runs a different leg of the discovery protocol.\n"
3394 )
3395 return 2
3396 if args.angles is not None and not args.finalize:
3397 sys.stderr.write(
3398 "[last30days] --angles only applies to --discover --finalize "
3399 "runs; add --finalize or drop the flag.\n"
3400 )
3401 return 2
3402 protocol_leg = (
3403 args.nominate_only or args.judgments is not None or args.finalize
3404 )
3405 if protocol_leg and args.mock and not args.save_dir:
3406 # Truthiness is right here: an empty --save-dir/env value means
3407 # "no save dir", and handoff state would land in the real config
3408 # dir - a side effect mock runs must never have.
3409 sys.stderr.write(
3410 "[last30days] mock protocol legs require --save-dir to stay "
3411 "side-effect-free: --mock with --nominate-only/--judgments/"
3412 "--finalize would otherwise write handoff state into the real "
3413 "config dir.\n"
3414 )
3415 return 2
3416 if protocol_leg:
3417 return _run_discover_protocol_leg(args, config)
3418 return _run_discover(args, config)
3419
3420 if args.discover_shallow:
3421 # Without --discover this flag would silently no-op into a full
3422 # research run - reject it instead of ignoring the requested mode.
3423 sys.stderr.write(
3424 "[last30days] --discover-shallow only applies to --discover runs; "
3425 "add --discover [domain] or drop the flag.\n"
3426 )
3427 return 2
3428
3429 # Same orphan rule for every protocol-leg flag: without --discover each
3430 # would silently no-op into a normal research run.
3431 for flag_label, present in (
3432 ("--nominate-only", args.nominate_only),
3433 ("--judgments", args.judgments is not None),
3434 ("--finalize", args.finalize),
3435 ):
3436 if present:
3437 sys.stderr.write(
3438 f"[last30days] {flag_label} only applies to --discover runs; "
3439 "add --discover [domain] or drop the flag.\n"
3440 )
3441 return 2
3442 if args.angles is not None:
3443 sys.stderr.write(
3444 "[last30days] --angles only applies to --discover --finalize runs; "
3445 "add --discover --finalize or drop the flag.\n"
3446 )
3447 return 2
3448
3449 if args.drill:
3450 if topic:
3451 sys.stderr.write(
3452 "[last30days] --drill uses the cached topic and cannot be "
3453 "combined with a new topic.\n"
3454 )
3455 return 2
3456 if args.publish_html and args.emit != "html":
3457 sys.stderr.write("[last30days] --publish-html requires --emit=html\n")
3458 return 2
3459 if args.dedicated_subreddits:
3460 config["_dedicated_subreddits"] = [
3461 value.strip().removeprefix("r/")
3462 for value in args.dedicated_subreddits.split(",")
3463 if value.strip()
3464 ]
3465 if args.polymarket_keywords:
3466 config["_polymarket_keywords"] = [
3467 value.strip().lower()
3468 for value in args.polymarket_keywords.split(",")
3469 if value.strip()
3470 ]
3471 return _run_drill(args, config)
3472
3473 if args.verify_freshness and not topic:
3474 return _run_cached_freshness(args, config)
3475
3476 if args.lookback_days is None:
3477 args.lookback_days = 30
3478
3479 if args.deep_research and not args.diagnose:
3480 from lib import planner as _planner
3481
3482 if not (
3483 config.get("PERPLEXITY_API_KEY")
3484 or config.get("OPENROUTER_API_KEY")
3485 ):
3486 print(
3487 "Error: --deep-research requires PERPLEXITY_API_KEY or "
3488 "OPENROUTER_API_KEY",
3489 file=sys.stderr,
3490 )
3491 return 1
3492 comparison_requested = any(
3493 value is not None
3494 for value in (
3495 args.competitors,
3496 args.competitors_list,
3497 args.competitors_plan,
3498 )
3499 ) or len(_planner._comparison_entities(topic, uncapped=True)) >= 2
3500 if comparison_requested:
3501 sys.stderr.write(
3502 "Error: --deep-research cannot be combined with competitor or vs-mode. "
3503 "It permits one paid Deep Research run per user action; run each topic "
3504 "separately.\n"
3505 )
3506 return 2
3507 config["_deep_research"] = True
3508 try:
3509 enable_deep_research_source(config)
3510 except ValueError as exc:
3511 print(f"Error: {exc}", file=sys.stderr)
3512 return 2
3513
3514 # Reject a misspelled configured register before remote submission or any
3515 # local source retrieval. Excluded modes resolve to default and remain
3516 # unaffected by the register setting.
3517 try:
3518 _audience_register_for_run(args, config, None)
3519 except ValueError as exc:
3520 sys.stderr.write(f"[last30days] {exc}\n")
3521 return 2
3522
3523 # Remote API path: when BOTH LAST30DAYS_API_KEY and LAST30DAYS_API_BASE are
3524 # set (and --mock is not), the search runs through the configured remote API
3525 # instead of local sources; no local provider keys are needed (see
3526 # lib/hosted.py). With either env var unset, behavior below is byte-identical
3527 # to local-only runs - there is no built-in endpoint.
3528 if (
3529 topic
3530 and resolved_corpus_dirs
3531 and env.read_secret_env("LAST30DAYS_API_KEY")
3532 and os.environ.get("LAST30DAYS_API_BASE")
3533 ):
3534 sys.stderr.write(
3535 "[last30days] Local corpus configured; bypassing the hosted backend so files stay on this machine.\n"
3536 )
3537 # An explicit --perplexity-search-type is per-invocation intent the hosted
3538 # backend cannot honor, so it runs locally, but only when a direct
3539 # PERPLEXITY_API_KEY can apply it; otherwise the switch would trade hosted
3540 # coverage for nothing. Key on the parsed CLI flag only: a value from
3541 # LAST30DAYS_PERPLEXITY_SEARCH_TYPE must never move routing.
3542 elif (
3543 topic
3544 and args.perplexity_search_type is not None
3545 and config.get("PERPLEXITY_API_KEY")
3546 and not args.diagnose
3547 and not args.mock
3548 and not args.record_fixtures
3549 and not args.deep_research
3550 and env.read_secret_env("LAST30DAYS_API_KEY")
3551 and os.environ.get("LAST30DAYS_API_BASE")
3552 ):
3553 sys.stderr.write(
3554 "[last30days] --perplexity-search-type set; bypassing the hosted backend "
3555 "because it does not apply the Perplexity search type.\n"
3556 )
3557 if (
3558 topic
3559 and not args.diagnose
3560 and not args.mock
3561 and not args.record_fixtures
3562 and env.read_secret_env("LAST30DAYS_API_KEY")
3563 and os.environ.get("LAST30DAYS_API_BASE")
3564 and not resolved_corpus_dirs
3565 and not args.deep_research
3566 and (args.perplexity_search_type is None or not config.get("PERPLEXITY_API_KEY"))
3567 ):
3568 if args.perplexity_search_type is not None:
3569 sys.stderr.write(
3570 "hosted backend does not apply --perplexity-search-type and no direct "
3571 "PERPLEXITY_API_KEY is configured to run it locally; skipping\n"
3572 )
3573 elif config.get("LAST30DAYS_PERPLEXITY_SEARCH_TYPE"):
3574 sys.stderr.write(
3575 "hosted backend does not apply LAST30DAYS_PERPLEXITY_SEARCH_TYPE; skipping\n"
3576 )
3577 if _freshness_enabled(args, config):
3578 if args.verify_freshness is True:
3579 sys.stderr.write(
3580 "[last30days] Freshness verification is not supported by the hosted backend; "
3581 "run locally or omit --verify-freshness.\n"
3582 )
3583 return 2
3584 sys.stderr.write(
3585 "hosted backend does not support freshness verification; skipping\n"
3586 )
3587 if args.emit == "json" and args.json_profile == "agent":
3588 sys.stderr.write(
3589 "[last30days] --json-profile=agent requires the local Report; "
3590 "the remote API backend only supports --json-profile=raw.\n"
3591 )
3592 return 2
3593 if args.x_posts is not None:
3594 # The envelope is a local-engine contract; the remote API has no
3595 # lane to receive it.
3596 sys.stderr.write(
3597 "[last30days] --x-posts is not supported by the hosted backend; "
3598 "run locally or omit --x-posts.\n"
3599 )
3600 return 2
3601 from lib import hosted
3602 depth = "deep" if args.deep else "quick" if args.quick else "default"
3603 try:
3604 audience = _audience_register_for_run(args, config, None)
3605 except ValueError as exc:
3606 sys.stderr.write(f"[last30days] {exc}\n")
3607 return 2
3608 hosted_kwargs = {
3609 "emit": args.emit,
3610 "save_dir": args.save_dir,
3611 "save_suffix": args.save_suffix or "",
3612 }
3613 if audience.name != "default":
3614 hosted_kwargs["register"] = audience.name
3615 return hosted.run_hosted(topic, depth, **hosted_kwargs)
3616
3617 requested_sources = resolve_requested_sources(args.search, config)
3618 if args.deep_research:
3619 requested_sources = add_deep_research_source(requested_sources)
3620 # Explicit --trustpilot-domain is user intent: activate the opt-in source
3621 # before diagnose/run so the flag cannot silently no-op (#873). Auto-resolve
3622 # hints are applied later and must not call this path.
3623 cli_trustpilot_domain = (
3624 args.trustpilot_domain.strip() if args.trustpilot_domain else ""
3625 )
3626 if cli_trustpilot_domain:
3627 requested_sources = activate_trustpilot_for_explicit_domain(
3628 config,
3629 requested_sources,
3630 reason=f"--trustpilot-domain={cli_trustpilot_domain}",
3631 )
3632 # Explicit --telegram-sources is user intent: activate the opt-in source
3633 # before diagnose/run so the flag cannot silently no-op (same pattern as
3634 # Trustpilot #873). Sets TELEGRAM_SOURCES in config for pipeline.
3635 cli_telegram_sources = (
3636 args.telegram_sources.strip() if args.telegram_sources else ""
3637 )
3638 if cli_telegram_sources:
3639 requested_sources = activate_telegram_for_explicit_sources(
3640 config,
3641 requested_sources,
3642 channels=cli_telegram_sources,
3643 )
3644 # Host-fetched X envelope: validated before diagnose so a present
3645 # envelope plans X in (available_sources) and a bad one fails closed here.
3646 x_posts_envelope: x_envelope.Envelope | None = None
3647 if args.x_posts is not None:
3648 if not topic:
3649 sys.stderr.write("[last30days] --x-posts requires a research topic.\n")
3650 return 2
3651 if _comparison_requested(args, topic):
3652 sys.stderr.write(
3653 "[last30days] --x-posts applies to a single-topic run; on a "
3654 "comparison run pass each entity's envelope through the "
3655 "x_posts field of its --competitors-plan entry.\n"
3656 )
3657 return 2
3658 try:
3659 x_posts_envelope = _read_x_envelope(
3660 args.x_posts, topic, args,
3661 x_handle=args.x_handle,
3662 x_related=args.x_related.split(",") if args.x_related else None,
3663 )
3664 except x_envelope.EnvelopeContractError as exc:
3665 sys.stderr.write(f"[last30days] {exc.message}\n")
3666 return 2
3667 diag = pipeline.diagnose(
3668 config, requested_sources, safe=args.diagnose,
3669 x_envelope=x_posts_envelope is not None,
3670 )
3671
3672 if args.diagnose:
3673 print(json.dumps(diag, indent=2, sort_keys=True))
3674 return 0
3675
3676 # Competitor sub-runs shallow-copy this config. The shared object makes the
3677 # paid Perplexity cap command-wide and thread-safe across that fanout. Keep
3678 # this runtime-only object out of the safe diagnose configuration contract.
3679 config["_perplexity_paid_budget"] = pipeline.PaidSourceBudget()
3680
3681 # Per-entity host-fetched X envelopes are validated here, on the main
3682 # thread and BEFORE the report-cache lookup, so a bad or stale one fails
3683 # closed (exit 2) instead of silently dropping that entity inside the
3684 # fan-out or being served from a cache built while it was still valid.
3685 comp_plan = parse_competitors_plan(args.competitors_plan)
3686 try:
3687 _attach_entity_envelopes(comp_plan, args)
3688 except x_envelope.EnvelopeContractError as exc:
3689 sys.stderr.write(f"[last30days] {exc.message}\n")
3690 return 2
3691
3692 if not topic:
3693 parser.print_usage(sys.stderr)
3694 return 2
3695 if args.publish_html and args.emit != "html":
3696 sys.stderr.write("[last30days] --publish-html requires --emit=html\n")
3697 return 2
3698
3699 synthesis_md = None
3700 if args.synthesis_file:
3701 if args.emit == "html":
3702 synthesis_md = read_synthesis_file(args.synthesis_file)
3703 else:
3704 sys.stderr.write("[last30days] Warning: --synthesis-file is only used with --emit=html; ignoring.\n")
3705
3706 if not os.environ.get("LAST30DAYS_SKIP_PREFLIGHT"):
3707 from lib import preflight
3708 refuse_msg = preflight.check_class_1_trap(topic)
3709 if refuse_msg:
3710 sys.stderr.write(refuse_msg)
3711 return 2
3712
3713 if (
3714 args.emit == "html"
3715 and synthesis_md is not None
3716 and not args.deep_research
3717 ):
3718 cached = _load_last_report_cache(
3719 topic,
3720 ttl_seconds=_report_cache_ttl_seconds(config),
3721 x_envelope_sha256=_x_envelope_digest(x_posts_envelope, comp_plan),
3722 )
3723 if cached is not None:
3724 cached_report, cached_entity_reports, cache_path = cached
3725 sys.stderr.write(
3726 f"[last30days] Reusing cached report data from {cache_path}\n"
3727 )
3728 sys.stderr.flush()
3729 if _freshness_enabled(args, config):
3730 _verify_report_set(
3731 cached_report,
3732 cached_entity_reports,
3733 allow_network=not args.mock,
3734 )
3735 _update_cached_freshness(
3736 cache_path,
3737 cached_report,
3738 cached_entity_reports,
3739 )
3740 return _render_save_and_print(
3741 args, cached_report, cached_entity_reports, synthesis_md, config
3742 )
3743 sys.stderr.write(
3744 "[last30days] No matching cached report data for "
3745 "--emit=html --synthesis-file; running fresh research.\n"
3746 )
3747 sys.stderr.flush()
3748
3749 progress = ui.ProgressDisplay(topic, show_banner=True)
3750 progress.start_processing()
3751
3752 depth = "deep" if args.deep else "quick" if args.quick else "default"
3753 # CLI overrides for the depth profile's result caps (issue #716). Stashed on
3754 # config so pipeline.run() can apply them without widening its signature; the
3755 # comparison path inherits them via `entity_config = dict(config)`.
3756 if args.max_results is not None:
3757 config["_max_results"] = args.max_results
3758 if args.max_per_source is not None:
3759 config["_max_per_source"] = args.max_per_source
3760 if args.max_source_fetches is not None:
3761 config["_max_source_fetches"] = args.max_source_fetches
3762 try:
3763 x_related = [h.strip() for h in args.x_related.split(",") if h.strip()] if args.x_related else None
3764 subreddits = [s.strip().removeprefix("r/") for s in args.subreddits.split(",") if s.strip()] if args.subreddits else None
3765 dedicated_subreddits = [s.strip().removeprefix("r/") for s in args.dedicated_subreddits.split(",") if s.strip()] if args.dedicated_subreddits else None
3766 tiktok_hashtags = [h.strip().lstrip("#") for h in args.tiktok_hashtags.split(",") if h.strip()] if args.tiktok_hashtags else None
3767 tiktok_creators = [c.strip().lstrip("@") for c in args.tiktok_creators.split(",") if c.strip()] if args.tiktok_creators else None
3768 ig_creators = [c.strip().lstrip("@") for c in args.ig_creators.split(",") if c.strip()] if args.ig_creators else None
3769 # Parse external plan if provided via --plan flag
3770 external_plan = None
3771 if args.plan:
3772 import json as _json
3773 plan_str = args.plan
3774 if os.path.isfile(plan_str):
3775 try:
3776 with open(plan_str, encoding="utf-8") as f:
3777 plan_str = f.read()
3778 except (OSError, UnicodeDecodeError) as exc:
3779 sys.stderr.write(f"[Planner] Cannot read --plan file: {exc}\n")
3780 raise SystemExit(2)
3781 try:
3782 external_plan = _json.loads(plan_str)
3783 except _json.JSONDecodeError as exc:
3784 sys.stderr.write(f"[Planner] Invalid --plan JSON: {exc}\n")
3785 # Fail fast instead of silently dropping to the internal planner
3786 # and burning a paid run the user did not ask for. Mirrors the
3787 # --plan file-read branch above and parse_competitors_plan.
3788 raise SystemExit(2)
3789 from lib import planner as _plan_validator
3790 try:
3791 _plan_validator.validate_external_plan(external_plan)
3792 except ValueError as exc:
3793 sys.stderr.write(f"[Planner] Invalid --plan schema: {exc}.\n")
3794 raise SystemExit(2)
3795
3796 # Auto-resolve: use web search to discover subreddits/handles before planning.
3797 # This is the engine-side equivalent of SKILL.md Steps 0.55/0.75 for platforms
3798 # without WebSearch (OpenClaw, Codex, raw CLI).
3799 repos_from_auto_resolve = False
3800 trustpilot_domain_is_hint = False
3801 # Resolve automatically for entity-shaped topics even without the flag.
3802 # A person or company topic whose handle the user did not supply is the
3803 # case where first-party evidence is hardest to protect: the handle is
3804 # absent from the topic and may never appear in retrieved mentions, so
3805 # nothing downstream can identify the subject's own posts. One web
3806 # search closes that. If it returns nothing, pipeline.run skips the X
3807 # relevance floor entirely — a noisier report beats losing evidence.
3808 # Skipped when a handle was already supplied, when an external plan
3809 # owns resolution, or in mock runs.
3810 if (
3811 not args.auto_resolve
3812 and not external_plan
3813 and not args.x_handle
3814 and not args.mock
3815 and _looks_like_entity_topic(topic)
3816 ):
3817 args.auto_resolve = True
3818 sys.stderr.write(
3819 "[AutoResolve] entity-shaped topic with no --x-handle; "
3820 "resolving the subject's handle so its own posts are not pruned\n"
3821 )
3822
3823 if args.auto_resolve and not external_plan:
3824 from lib import resolve
3825 resolution = resolve.auto_resolve(topic, config)
3826 if resolution.get("subreddits") and not subreddits:
3827 subreddits = resolution["subreddits"]
3828 sys.stderr.write(f"[AutoResolve] Subreddits: {', '.join(subreddits)}\n")
3829 if resolution.get("x_handle") and not args.x_handle:
3830 args.x_handle = resolution["x_handle"]
3831 sys.stderr.write(f"[AutoResolve] X handle: @{args.x_handle}\n")
3832 # Empty x_handle is intentional: do not invent a lexical stand-in.
3833 # pipeline.run treats an unidentified subject as "skip the X floor".
3834 if resolution.get("github_user") and not args.github_user:
3835 args.github_user = resolution["github_user"]
3836 sys.stderr.write(f"[AutoResolve] GitHub user: @{args.github_user}\n")
3837 if resolution.get("github_repos") and not args.github_repo:
3838 args.github_repo = ",".join(resolution["github_repos"])
3839 # auto_resolve already canonicalized via canonicalize_github_repos(cap=5);
3840 # mark so we don't re-canonicalize below and clobber its relevance order.
3841 repos_from_auto_resolve = True
3842 sys.stderr.write(f"[AutoResolve] GitHub repos: {args.github_repo}\n")
3843 if resolution.get("trustpilot_domain") and not args.trustpilot_domain:
3844 # Hint provenance matters: only user-set flags are verbatim-final;
3845 # a resolved hint retries via the CLI search when it misses.
3846 args.trustpilot_domain = resolution["trustpilot_domain"]
3847 trustpilot_domain_is_hint = True
3848 sys.stderr.write(f"[AutoResolve] Trustpilot domain: {args.trustpilot_domain} (hint)\n")
3849 if resolution.get("context"):
3850 # Inject context into external_plan metadata for the planner to use
3851 if not external_plan:
3852 external_plan = None # planner will use its own, but with context
3853 # Store context for the planner prompt injection
3854 config["_auto_resolve_context"] = resolution["context"]
3855 sys.stderr.write(f"[AutoResolve] Context: {resolution['context'][:80]}...\n")
3856
3857 github_user = args.github_user.lstrip("@").lower() if args.github_user else None
3858 github_repos = [r.strip() for r in args.github_repo.split(",") if r.strip() and "/" in r.strip()] if args.github_repo else None
3859 trustpilot_domain = args.trustpilot_domain.strip() if args.trustpilot_domain else None
3860
3861 comp_enabled, comp_count, comp_explicit = resolve_competitors_args(args)
3862 # comp_plan was parsed, and its per-entity envelopes validated, before
3863 # the report-cache lookup above.
3864
3865 # Plan-level trustpilot_domain pins are the same user intent as the CLI
3866 # flag (already activated above). Auto-resolve hints must not activate.
3867 if plan_has_explicit_trustpilot_domain(comp_plan):
3868 requested_sources = activate_trustpilot_for_explicit_domain(
3869 config,
3870 requested_sources,
3871 reason="competitors-plan trustpilot_domain",
3872 )
3873
3874 # Only canonicalize when repos came from a user-supplied --github-repo flag.
3875 # When repos_from_auto_resolve is True, auto_resolve already ran
3876 # canonicalize_github_repos(cap=5) and ranked by relevance; re-running here
3877 # with cap=None can re-sort by topic-slug match and lose that ordering.
3878 if github_repos and not repos_from_auto_resolve:
3879 from lib import resolve as resolve_lib
3880 original_github_repos = github_repos[:]
3881 github_repos = resolve_lib.canonicalize_github_repos(topic, github_repos, cap=None)
3882 if github_repos != original_github_repos:
3883 sys.stderr.write(
3884 "[GitHub] Canonicalized repos: "
3885 f"{','.join(original_github_repos)} -> {','.join(github_repos)}\n"
3886 )
3887
3888 # Polymarket disambiguation: if user passed --polymarket-keywords,
3889 # store on config so the polymarket adapter can filter matches.
3890 if args.polymarket_keywords:
3891 keywords = [
3892 k.strip().lower()
3893 for k in args.polymarket_keywords.split(",")
3894 if k.strip()
3895 ]
3896 if keywords:
3897 config["_polymarket_keywords"] = keywords
3898
3899 # Product keyword for the amazon source. Carried on config rather than
3900 # threaded through the run signature (the _polymarket_keywords idiom):
3901 # it is one optional string consumed in exactly two places.
3902 if getattr(args, "amazon_query", None):
3903 config["_amazon_query"] = args.amazon_query.strip()
3904 # Unlike --trustpilot-domain, this flag deliberately does NOT
3905 # auto-activate its source: the lane spends metered credits, so
3906 # turning it on stays an explicit request. But silence is the
3907 # wrong failure mode -- a model that resolves the keyword and
3908 # forgets the --search token would otherwise get no signal at
3909 # all that the flag did nothing.
3910 _amazon_requested = (
3911 (requested_sources and "amazon" in requested_sources)
3912 or "amazon" in str(config.get("INCLUDE_SOURCES") or "").lower()
3913 )
3914 if not _amazon_requested:
3915 sys.stderr.write(
3916 "[Amazon] --amazon-query was set but the amazon source was not "
3917 "requested; add it to --search (e.g. --search reddit,x,amazon) "
3918 "or set INCLUDE_SOURCES=amazon. Ignoring the keyword.\n"
3919 )
3920
3921 # Advertiser page override for the meta_ads source. Same shape as
3922 # --amazon-query (config-carried, warn-not-activate) and for the same
3923 # reason: the lane spends metered credits per call.
3924 if getattr(args, "meta_ads_page", None):
3925 page_id = parse_meta_ads_page(args.meta_ads_page)
3926 if not page_id:
3927 sys.stderr.write(
3928 "[Meta Ads] --meta-ads-page must be a numeric Ad Library page id "
3929 "or an Ad Library URL containing view_all_page_id; a facebook.com "
3930 "vanity URL is not a page id. Ignoring the override.\n"
3931 )
3932 else:
3933 config["_meta_ads_page"] = page_id
3934 _meta_ads_requested = (
3935 (requested_sources and "meta_ads" in requested_sources)
3936 or "meta_ads" in str(config.get("INCLUDE_SOURCES") or "").lower()
3937 )
3938 if not _meta_ads_requested:
3939 sys.stderr.write(
3940 "[Meta Ads] --meta-ads-page was set but the meta_ads source "
3941 "was not requested; add it to --search (e.g. --search "
3942 "reddit,x,meta_ads) or set INCLUDE_SOURCES=meta_ads. "
3943 "Ignoring the page.\n"
3944 )
3945
3946 # vs-mode / plan routing: split a vs-topic into main + peers unless
3947 # discover-N or an explicit --competitors-list already decided who runs.
3948 topic, comp_enabled, comp_count, comp_explicit = apply_vs_competitor_routing(
3949 topic,
3950 competitors_flag=args.competitors,
3951 comp_enabled=comp_enabled,
3952 comp_count=comp_count,
3953 comp_explicit=comp_explicit,
3954 comp_plan=comp_plan,
3955 )
3956 if comp_enabled:
3957 config["_perplexity_paid_budget"] = pipeline.PaidSourceBudget(
3958 owner=topic,
3959 )
3960
3961 # Plan alone with zero peers (empty/invalid JSON object, or all entries
3962 # skipped) must not fall through to discover-N with a misleading abort.
3963 if (
3964 comp_enabled
3965 and not comp_explicit
3966 and args.competitors is None
3967 and args.competitors_plan
3968 ):
3969 sys.stderr.write(
3970 "[Competitors] --competitors-plan has no usable peer entries "
3971 "(and the topic is not a vs-comparison). Pass a non-empty plan, "
3972 "a vs-topic, --competitors-list, or --competitors N.\n"
3973 )
3974 return 2
3975
3976 # Dedicated subs ride the config dict (already threaded to every source
3977 # fetch) so the keyless Reddit path can pull them floor-exempt without
3978 # widening pipeline.run / _retrieve_stream signatures.
3979 if dedicated_subreddits:
3980 config["_dedicated_subreddits"] = dedicated_subreddits
3981
3982 def _main_runner() -> schema.Report:
3983 r = pipeline.run(
3984 topic=topic,
3985 config=config,
3986 depth=depth,
3987 requested_sources=requested_sources,
3988 mock=args.mock,
3989 x_handle=args.x_handle,
3990 x_related=x_related,
3991 web_backend=args.web_backend,
3992 external_plan=external_plan,
3993 subreddits=subreddits,
3994 tiktok_hashtags=tiktok_hashtags,
3995 tiktok_creators=tiktok_creators,
3996 ig_creators=ig_creators,
3997 lookback_days=args.lookback_days,
3998 as_of_date=args.as_of_date,
3999 github_user=github_user,
4000 github_repos=github_repos,
4001 trustpilot_domain=trustpilot_domain,
4002 trustpilot_domain_is_hint=trustpilot_domain_is_hint,
4003 internal_subrun=comp_enabled,
4004 hiring_signals_mode=args.hiring_signals,
4005 save_dir=args.save_dir,
4006 corpus_dirs=args.corpus,
4007 corpus_all_time=args.corpus_all_time,
4008 x_posts=(
4009 comp_plan.get(topic.strip().lower(), {}).get("_x_envelope")
4010 if comp_enabled else x_posts_envelope
4011 ),
4012 )
4013 r.artifacts["resolved"] = {
4014 "entity": topic,
4015 "x_handle": (args.x_handle or "").lstrip("@"),
4016 "subreddits": list(subreddits or []),
4017 "github_user": (github_user or ""),
4018 "github_repos": list(github_repos or []),
4019 "trustpilot_domain": (trustpilot_domain or ""),
4020 "context": config.get("_auto_resolve_context", "") or "",
4021 }
4022 return r
4023
4024 if comp_enabled:
4025 from lib import competitors as competitors_mod
4026 from lib import fanout, resolve as resolve_mod
4027
4028 if comp_explicit:
4029 discovered = comp_explicit
4030 else:
4031 if not resolve_mod._has_backend(config) and not args.mock:
4032 sys.stderr.write(
4033 "[Competitors] Cannot auto-discover peers without help.\n"
4034 "\n"
4035 "RECOMMENDED PATH (hosting reasoning models — Claude Code, Codex, "
4036 "Hermes, Gemini, any agent with a WebSearch tool): YOU have "
4037 "WebSearch. Use it to run full Step 0.55 per entity, then invoke "
4038 "the engine with a vs-topic plus --competitors-plan:\n"
4039 " 1. WebSearch for '{topic} competitors' or '{topic} alternatives'.\n"
4040 " 2. For each peer, WebSearch for handles/subs/github (Step 0.55).\n"
4041 " 3. Re-invoke: /last30days '{topic} vs {peer1} vs {peer2}' "
4042 "--competitors-plan '{\"Peer1\":{\"x_handle\":\"h1\",\"subreddits\":"
4043 "[\"s1\"],...},\"Peer2\":{...}}'.\n"
4044 "See SKILL.md 'Competitor mode' for the full protocol.\n"
4045 "\n"
4046 "HEADLESS / CRON PATH (no hosting model available): set "
4047 "BRAVE_API_KEY / EXA_API_KEY / SERPER_API_KEY / PARALLEL_API_KEY / "
4048 "PERPLEXITY_API_KEY / OPENROUTER_API_KEY and re-run.\n"
4049 "\n"
4050 "MINIMUM ESCAPE HATCH: pass --competitors-list 'A,B,C' to skip "
4051 "discovery. Without --competitors-plan, peer sub-runs fall back to "
4052 "planner defaults and produce visibly thinner data than the main.\n"
4053 )
4054 return 2
4055 discovered = competitors_mod.discover_competitors(
4056 topic, comp_count, config, lookback_days=args.lookback_days,
4057 )
4058 if not discovered:
4059 sys.stderr.write(
4060 f"[Competitors] No peers discovered for {topic!r}; aborting "
4061 "comparison run. Pass --competitors-list to override.\n"
4062 )
4063 return 2
4064
4065 # run_competitor_fanout keys its results by label, so two
4066 # submissions sharing one collapse to a single report while the
4067 # returned list still carries two entries. That yields a
4068 # comparison of an entity against itself, and it hides a failed
4069 # main topic from the survivor check below: the duplicate peer's
4070 # report answers for the label the main run was supposed to fill.
4071 distinct_peers: list[str] = []
4072 claimed_labels = {comparison_label_key(topic)}
4073 for peer in discovered:
4074 key = comparison_label_key(peer)
4075 if key in claimed_labels:
4076 sys.stderr.write(
4077 f"[Competitors] Dropping {peer!r}: duplicates the main "
4078 "topic or an earlier peer.\n"
4079 )
4080 continue
4081 claimed_labels.add(key)
4082 distinct_peers.append(peer)
4083 if not distinct_peers:
4084 sys.stderr.write(
4085 f"[Competitors] No peer distinct from {topic!r} remains; "
4086 "there is nothing to compare against. Pass "
4087 "--competitors-list with distinct entities.\n"
4088 )
4089 return 2
4090 discovered = distinct_peers
4091
4092 sys.stderr.write(
4093 f"[Competitors] Comparing: {topic} vs " + " vs ".join(discovered) + "\n"
4094 )
4095
4096 def _competitor_runner(entity: str) -> schema.Report:
4097 # Deep-copy config so per-entity auto_resolve context does not
4098 # leak across sub-runs. Each sub-run writes its own
4099 # `_auto_resolve_context` into its local config copy.
4100 entity_config = dict(config)
4101 # The Amazon keyword is entity-SPECIFIC, unlike the depth caps
4102 # this shallow copy exists to inherit. Leaving the main topic's
4103 # keyword in place would search Weber SKUs for a Traeger peer,
4104 # render a rival's products as that peer's buyer evidence, and
4105 # multiply the metered spend by the number of entities. Drop it
4106 # so each peer derives its own keyword from its own topic; a
4107 # per-entity keyword can ride in the --competitors-plan entry.
4108 entity_config.pop("_amazon_query", None)
4109 # An advertiser page is per-entity state by definition: left in
4110 # place it would render one brand's ads as every peer's.
4111 entity_config.pop("_meta_ads_page", None)
4112 plan_entry = comp_plan.get(entity.strip().lower(), {})
4113 resolved = {
4114 "entity": entity,
4115 "x_handle": "",
4116 "subreddits": [],
4117 "github_user": "",
4118 "github_repos": [],
4119 "trustpilot_domain": "",
4120 "context": "",
4121 }
4122 # Skip engine-internal auto_resolve when the hosting model
4123 # pre-resolved via --competitors-plan (saves a redundant
4124 # round-trip and makes per-entity Step 0.55 purely
4125 # hosting-model-driven).
4126 plan_covers_fully = bool(plan_entry.get("x_handle")) and bool(
4127 plan_entry.get("subreddits")
4128 )
4129 if (
4130 not args.mock
4131 and not plan_covers_fully
4132 and resolve_mod._has_backend(entity_config)
4133 ):
4134 try:
4135 r = resolve_mod.auto_resolve(entity, entity_config)
4136 except Exception as exc:
4137 sys.stderr.write(
4138 f"[Competitors] auto_resolve failed for {entity!r}: "
4139 f"{type(exc).__name__}: {exc}\n"
4140 )
4141 r = {}
4142 resolved["x_handle"] = r.get("x_handle", "") or ""
4143 resolved["subreddits"] = list(r.get("subreddits") or [])
4144 resolved["github_user"] = r.get("github_user", "") or ""
4145 resolved["github_repos"] = list(r.get("github_repos") or [])
4146 resolved["trustpilot_domain"] = r.get("trustpilot_domain", "") or ""
4147 resolved["context"] = r.get("context", "") or ""
4148 kwargs = subrun_kwargs_for(entity, plan_entry, resolved=resolved)
4149 # Record effective per-entity targeting for the Resolved block.
4150 resolved_effective = {
4151 "entity": entity,
4152 "x_handle": kwargs["x_handle"] or "",
4153 "subreddits": kwargs["subreddits"] or [],
4154 "github_user": kwargs["github_user"] or "",
4155 "github_repos": kwargs["github_repos"] or [],
4156 "trustpilot_domain": kwargs["trustpilot_domain"] or "",
4157 "context": kwargs["_context"],
4158 }
4159 if kwargs["_context"]:
4160 entity_config["_auto_resolve_context"] = kwargs["_context"]
4161 sys.stderr.write(
4162 f"[Competitors] {entity}: "
4163 f"x=@{resolved_effective['x_handle'] or '-'} "
4164 f"subs={len(resolved_effective['subreddits'])} "
4165 f"gh={resolved_effective['github_user'] or '-'} "
4166 f"({'plan' if plan_entry else 'auto'})\n"
4167 )
4168 report = pipeline.run(
4169 topic=entity,
4170 config=entity_config,
4171 depth=depth,
4172 requested_sources=requested_sources,
4173 mock=args.mock,
4174 x_handle=kwargs["x_handle"],
4175 x_related=kwargs["x_related"],
4176 subreddits=kwargs["subreddits"],
4177 github_user=kwargs["github_user"],
4178 github_repos=kwargs["github_repos"],
4179 trustpilot_domain=kwargs["trustpilot_domain"],
4180 trustpilot_domain_is_hint=kwargs["_trustpilot_domain_is_hint"],
4181 web_backend=args.web_backend,
4182 lookback_days=args.lookback_days,
4183 as_of_date=args.as_of_date,
4184 hiring_signals_mode=args.hiring_signals,
4185 internal_subrun=True,
4186 save_dir=args.save_dir,
4187 corpus_dirs=args.corpus,
4188 corpus_all_time=args.corpus_all_time,
4189 x_posts=plan_entry.get("_x_envelope"),
4190 )
4191 report.artifacts["resolved"] = resolved_effective
4192 return report
4193
4194 entity_reports = fanout.run_competitor_fanout(
4195 main_topic=topic,
4196 main_runner=_main_runner,
4197 competitors=discovered,
4198 competitor_runner=_competitor_runner,
4199 )
4200 # run_competitor_fanout drops a failed sub-run from the list, and
4201 # the render takes entity_reports[0] as the comparison's subject.
4202 # Without this check, a main topic that raised while >=2 peers
4203 # succeeded silently promoted a competitor to be the subject: the
4204 # report was headed by that peer, saved under its slug, and the
4205 # topic the user actually asked about went unmentioned.
4206 survived = {label for label, _ in entity_reports}
4207 dropped = [
4208 label for label in (topic, *discovered) if label not in survived
4209 ]
4210 if topic not in survived:
4211 progress.end_processing()
4212 sys.stderr.write(
4213 f"[Competitors] The main topic {topic!r} failed; "
4214 f"{len(entity_reports)} competitor sub-run(s) survived. "
4215 "Refusing to render a comparison headed by a competitor. "
4216 "Check the warnings above.\n"
4217 )
4218 return 1
4219 if len(entity_reports) < 2:
4220 progress.end_processing()
4221 sys.stderr.write(
4222 f"[Competitors] Fewer than 2 sub-runs survived ({len(entity_reports)}); "
4223 "cannot render a comparison. Re-run without --competitors or check the "
4224 "warnings above.\n"
4225 )
4226 return 1
4227 report = entity_reports[0][1]
4228 if dropped:
4229 # A narrower comparison than the user asked for is a result
4230 # they need to see, not a silent substitution.
4231 report.warnings.append(
4232 "Comparison is incomplete: "
4233 f"{len(dropped)} of {len(discovered) + 1} entities failed and "
4234 f"were dropped ({', '.join(dropped)})."
4235 )
4236 else:
4237 entity_reports = None
4238 report = _main_runner()
4239 except Exception as exc:
4240 progress.end_processing()
4241 progress.show_error(str(exc))
4242 raise
4243 if _freshness_enabled(args, config):
4244 _verify_report_set(report, entity_reports, allow_network=not args.mock)
4245
4246 _show_runtime_ui(
4247 report, progress, diag,
4248 suppress_web_promo=bool(external_plan or comp_plan),
4249 )
4250 _write_last_run(
4251 original_topic, report, entity_reports=entity_reports,
4252 x_envelope_sha256=_x_envelope_digest(x_posts_envelope, comp_plan),
4253 )
4254 # LAST30DAYS_STORE env var = persistence default-on. Read both os.environ
4255 # (for shell-exported users) and config (for users who set it in
4256 # ~/.config/last30days/.env, which env.py loads but does not propagate
4257 # to os.environ). Mirrors the LAST30DAYS_DEBUG / LAST30DAYS_SKIP_PREFLIGHT
4258 # convention; env-var or config wins, with `--store` flag still working.
4259 _store_env = (
4260 os.environ.get("LAST30DAYS_STORE")
4261 or config.get("LAST30DAYS_STORE")
4262 or ""
4263 ).lower()
4264 if args.store or _store_env in ("1", "true", "yes"):
4265 counts = persist_report(report, store_db=_scoped_store_db(args))
4266 sys.stderr.write(
4267 f"[last30days] Stored {counts['new']} new, {counts['updated']} updated findings\n"
4268 )
4269 sys.stderr.flush()
4270
4271 # Show quality nudge if applicable. Explicit hiring-signal runs are
4272 # intentionally jobs-focused, so generic source setup advice is noise.
4273 if not args.hiring_signals:
4274 try:
4275 from lib import quality_nudge
4276 from lib import youtube_yt as _youtube_yt
4277 # Populate transcript-fetch ratio so quality_nudge can detect the
4278 # degraded-YouTube failure mode (videos returned but transcripts
4279 # silently failed - typically a stale yt-dlp binary).
4280 youtube_items = report.items_by_source.get("youtube") or []
4281 _yt_fetch_stats = _youtube_yt.get_transcript_fetch_stats()
4282 instagram_items = report.items_by_source.get("instagram") or []
4283 research_results = {
4284 "active_sources": diag.get("available_sources") or [],
4285 "youtube_videos_count": len(youtube_items),
4286 "youtube_transcripts_count": sum(
4287 1 for it in youtube_items
4288 if (it.metadata.get("transcript_highlights") or it.metadata.get("transcript_snippet"))
4289 ),
4290 "youtube_error": report.errors_by_source.get("youtube"),
4291 "x_error": report.errors_by_source.get("x"),
4292 # Captions-disabled videos can never produce a transcript regardless
4293 # of yt-dlp version; subtract them from the degraded-ratio
4294 # denominator so a single uploader-disabled video does not trip the
4295 # "stale yt-dlp" nudge.
4296 "youtube_captions_disabled_count": sum(
4297 1 for it in youtube_items if it.metadata.get("captions_disabled")
4298 ),
4299 # Actual yt-dlp fetch outcomes for this run. The counts above are
4300 # computed from post-pruning items, so they can't tell "fetches
4301 # failed (stale binary)" from "fetches succeeded but the videos
4302 # were pruned downstream"; the latter was producing false
4303 # stale-yt-dlp nudges (#531).
4304 "youtube_transcript_fetch_attempts": _yt_fetch_stats["attempts"],
4305 "youtube_transcript_fetch_failures": _yt_fetch_stats["failures"],
4306 # Track Instagram returned-zero-items so quality_nudge can detect
4307 # the silent-failure case (SC configured but the v2 reels endpoint
4308 # 500'd through both the original query and the hashtag retry).
4309 "instagram_items_count": len(instagram_items),
4310 }
4311 quality = quality_nudge.compute_quality_score(config, research_results)
4312 if quality.get("nudge_text"):
4313 sys.stderr.write(f"\n{quality['nudge_text']}\n")
4314 sys.stderr.flush()
4315 except Exception:
4316 pass
4317
4318 # Signal to render_compact whether pre-research flags were supplied.
4319 # Used to emit a Pre-Research Status warning when the model skipped
4320 # Step 0.5 / 0.55 and invoked the engine bare on an eligible topic.
4321 pre_research_flags_present = bool(
4322 args.x_handle
4323 or args.github_user
4324 or args.subreddits
4325 or args.plan
4326 or args.auto_resolve
4327 or args.tiktok_creators
4328 or args.ig_creators
4329 )
4330 report.artifacts["pre_research_flags_present"] = pre_research_flags_present
4331
4332 exit_code = _render_save_and_print(args, report, entity_reports, synthesis_md, config)
4333 if args.emit in {"compact", "md", "brief"}:
4334 x_omission = _optional_x_omission_text(diag, requested_sources)
4335 if x_omission:
4336 sys.stderr.write(f"\n{x_omission}\n")
4337 sys.stderr.flush()
4338 return exit_code
4339
4340
4341 if __name__ == "__main__":
4342 _install_sigterm_handler()
4343 raise SystemExit(main())
4344
4344 lines PYTHON