#!/usr/bin/env python3 """ archivebox run [--daemon] [--crawl-id=...] [--snapshot-id=...] [--binary-id=...] Unified command for processing queued work on the shared abx-dl bus. Modes: - With stdin JSONL: Process piped records, exit when complete - Without stdin (TTY): Run the background runner in foreground until killed - --crawl-id: Run the crawl runner for a specific crawl only - --snapshot-id: Run a specific snapshot through its parent crawl - --binary-id: Emit a BinaryRequestEvent for a specific Binary row Examples: # Run the background runner in foreground archivebox run # Run as daemon (don't exit on idle) archivebox run --daemon # Process specific records (pipe any JSONL type, exits when done) archivebox snapshot list --status=queued | archivebox run archivebox archiveresult list --status=failed | archivebox run archivebox crawl list --status=queued | archivebox run # Mixed types work too cat mixed_records.jsonl | archivebox run # Run the crawl runner for a specific crawl archivebox run --crawl-id=019b7e90-04d0-73ed-adec-aad9cfcd863e # Run one snapshot from an existing crawl archivebox run --snapshot-id=019b7e90-5a8e-712c-9877-2c70eebe80ad # Run one queued binary install directly on the bus archivebox run --binary-id=019b7e90-5a8e-712c-9877-2c70eebe80ad """ __package__ = "archivebox.cli" __command__ = "archivebox run" import asyncio import os import signal import sys from collections import defaultdict import rich_click as click from rich import print as rprint RUNNER_DAEMON_ENV = "ARCHIVEBOX_RUNNER_DAEMON" def _exit_daemon_runner_on_signal(sig: signal.Signals) -> None: # A supervised `archivebox run --daemon` is intentionally a disposable # child. If it receives SIGINT/SIGTERM directly, exit with the conventional # signal status so supervisord treats it as an unexpected worker death and # restarts only the runner. The parent `archivebox server` owns supervisord # shutdown and must not be pulled down by a killed daemon worker. os._exit(128 + int(sig)) def process_stdin_records() -> int: """ Process JSONL records from stdin. Create-or-update behavior: - Records WITHOUT id: Create via Model.from_json(), then queue - Records WITH id: Lookup existing, re-queue for processing Outputs JSONL of all processed records (for chaining). Handles Crawl and Snapshot work records. ArchiveResult records are accepted as references to their parent Snapshot and plugin. Returns exit code (0 = success, 1 = error). """ from django.utils import timezone from archivebox.misc.jsonl import ( read_stdin, write_record, TYPE_CRAWL, TYPE_SNAPSHOT, TYPE_ARCHIVERESULT, TYPE_BINARYREQUEST, TYPE_BINARY, ) from archivebox.base_models.models import get_or_create_system_user_pk from archivebox.core.models import Snapshot, ArchiveResult from archivebox.api.v1_core import _uuid_ref_query from archivebox.crawls.models import Crawl from archivebox.core.shutdown_util import foreground_parent_watchdog, foreground_shutdown_signals from archivebox.machine.models import Binary from archivebox.services.runner import run_binary, run_crawl records = list(read_stdin()) is_tty = sys.stdout.isatty() if not records: return 0 # Nothing to process created_by_id = get_or_create_system_user_pk() queued_count = 0 output_records = [] full_crawl_ids: set[str] = set() snapshot_ids_by_crawl: dict[str, set[str]] = defaultdict(set) plugin_names_by_crawl: dict[str, set[str]] = defaultdict(set) run_all_plugins_for_crawl: set[str] = set() binary_ids: list[str] = [] for record in records: record_type = record.get("type", "") record_id = record.get("id") try: if record_type == TYPE_CRAWL: if record_id: # Existing crawl - re-queue try: crawl = Crawl.objects.get(id=record_id) except Crawl.DoesNotExist: crawl = Crawl.from_json(record, overrides={"created_by_id": created_by_id}) else: # New crawl - create it crawl = Crawl.from_json(record, overrides={"created_by_id": created_by_id}) if crawl: crawl.update_and_requeue( status=Crawl.StatusChoices.QUEUED, retry_at=timezone.now(), ) full_crawl_ids.add(str(crawl.id)) run_all_plugins_for_crawl.add(str(crawl.id)) output_records.append(crawl.to_json()) queued_count += 1 elif record_type == TYPE_SNAPSHOT or (record.get("url") and not record_type): if record_id: # Existing snapshot - re-queue try: snapshot = Snapshot.objects.get(id=record_id) except Snapshot.DoesNotExist: snapshot = Snapshot.from_json(record, overrides={"created_by_id": created_by_id}) else: # New snapshot - create it snapshot = Snapshot.from_json(record, overrides={"created_by_id": created_by_id}) if snapshot: snapshot.queue_for_extraction() crawl_id = str(snapshot.crawl_id) snapshot_ids_by_crawl[crawl_id].add(str(snapshot.id)) run_all_plugins_for_crawl.add(crawl_id) output_records.append(snapshot.to_json()) queued_count += 1 elif record_type == TYPE_ARCHIVERESULT: snapshot_id = str(record.get("snapshot_id") or "") plugin_name = str(record.get("plugin") or "") archiveresult = None if not snapshot_id and record_id: archiveresult = ArchiveResult.objects.filter(_uuid_ref_query("id", str(record_id))).select_related("snapshot").first() if archiveresult: snapshot_id = str(archiveresult.snapshot_id) plugin_name = plugin_name or archiveresult.plugin snapshot = Snapshot.objects.filter(id=snapshot_id).first() if snapshot_id else None if snapshot: snapshot.queue_for_extraction() crawl_id = str(snapshot.crawl_id) snapshot_ids_by_crawl[crawl_id].add(str(snapshot.id)) if plugin_name: plugin_names_by_crawl[crawl_id].add(str(plugin_name)) output_records.append(archiveresult.to_json() if archiveresult else record) queued_count += 1 elif record_type in {TYPE_BINARYREQUEST, TYPE_BINARY}: if record_id: try: binary = Binary.objects.get(id=record_id) except Binary.DoesNotExist: binary = Binary.from_json(record) else: binary = Binary.from_json(record) if binary: binary.retry_at = timezone.now() if binary.status != Binary.StatusChoices.INSTALLED: binary.status = Binary.StatusChoices.QUEUED binary.save() binary_ids.append(str(binary.id)) output_records.append(binary.to_json()) queued_count += 1 else: # Unknown type - pass through output_records.append(record) except Exception as e: rprint(f"[yellow]Error processing record: {e}[/yellow]", file=sys.stderr) continue # Output all processed records (for chaining) if not is_tty: for rec in output_records: write_record(rec) if queued_count == 0: rprint("[yellow]No records to process[/yellow]", file=sys.stderr) return 0 rprint(f"[blue]Processing {queued_count} records...[/blue]", file=sys.stderr) for binary_id in binary_ids: run_binary(binary_id) targeted_crawl_ids = full_crawl_ids | set(snapshot_ids_by_crawl) if targeted_crawl_ids: for crawl_id in sorted(targeted_crawl_ids): try: crawl = Crawl.objects.get(id=crawl_id) except Crawl.DoesNotExist: continue if not crawl.claim_processing_lock(lock_seconds=10): rprint(f"[yellow]Crawl {crawl_id} is already owned by another runner[/yellow]", file=sys.stderr) return 1 with foreground_shutdown_signals(), foreground_parent_watchdog(): run_crawl( crawl_id, snapshot_ids=None if crawl_id in full_crawl_ids else sorted(snapshot_ids_by_crawl[crawl_id]), selected_plugins=None if crawl_id in run_all_plugins_for_crawl else sorted(plugin_names_by_crawl[crawl_id]), ) return 0 def run_runner( daemon: bool = False, crawl_id: str | None = None, maintenance_only: bool = False, ) -> int: """ Run the background runner loop. Args: daemon: Run forever (don't exit when idle) Returns exit code (0 = success, 1 = error). """ from archivebox.config import CONSTANTS from archivebox.core.shutdown_util import foreground_parent_watchdog, foreground_shutdown_signals from archivebox.machine.models import Machine, Process from archivebox.core.takeover_util import enter_single_runner_gate, standby_until_foreground_runner_needed from archivebox.core.recovery_util import recover_orchestrator_state from archivebox.services.runner import run_pending_crawls Machine.current() current = Process.current() root_command = current.root if daemon and root_command.process_type in ( Process.TypeChoices.SERVER, Process.TypeChoices.ADD, Process.TypeChoices.UPDATE, ): # Server-owned daemon runners are persistent supervisor workers, but # foreground add/update commands are allowed to borrow runner/sonic # leadership without taking down Daphne. Waiting here keeps the worker # on the normal runner path while preventing a server restart loop from # immediately stealing the single-runner gate back from the newer CLI. standby_until_foreground_runner_needed(root_command, data_dir=CONSTANTS.DATA_DIR) if not enter_single_runner_gate(current, data_dir=CONSTANTS.DATA_DIR): current.mark_exited() return 0 recover_orchestrator_state(include_chrome=crawl_id is None, crawl_id=crawl_id) if crawl_id: from django.utils import timezone from archivebox.crawls.models import Crawl crawl = Crawl.objects.filter(id=crawl_id, status__in=Crawl.RUNNABLE_STATES).first() now = timezone.now() # Winning the single-runner gate terminates and waits for every older # runner. Requeue the explicitly requested crawl so an abandoned # future ownership lease cannot hide it from the unified scheduler. if crawl is not None: crawl.update_and_requeue(retry_at=now, refresh=False) # Only a foreground `archivebox add` gets the interactive "abort current # hook, continue/retry, second Ctrl+C exits" flow. Server/update/run owned # orchestrators should shut down immediately and cleanly on the first signal. interactive_interrupts = current.root.process_type == Process.TypeChoices.ADD if daemon: os.environ[RUNNER_DAEMON_ENV] = "1" try: with ( foreground_shutdown_signals( on_signal=_exit_daemon_runner_on_signal if daemon else None, raise_on_first_signal=not daemon, ), foreground_parent_watchdog(enabled=not daemon), ): run_pending_crawls( daemon=daemon, crawl_id=crawl_id, maintenance_only=maintenance_only, interactive_interrupts=interactive_interrupts, ) return 0 except KeyboardInterrupt: return 0 except asyncio.CancelledError as e: if daemon: rprint(f"[red]Runner cancelled unexpectedly: {type(e).__name__}: {e}[/red]", file=sys.stderr) return 1 return 0 except Exception as e: rprint(f"[red]Runner error: {type(e).__name__}: {e}[/red]", file=sys.stderr) return 1 finally: current.refresh_from_db() if current.status != Process.StatusChoices.EXITED: current.mark_exited() @click.command() @click.option("--daemon", "-d", is_flag=True, help="Run forever (don't exit on idle)") @click.option("--crawl-id", help="Run the crawl runner for a specific crawl only") @click.option("--snapshot-id", help="Run one snapshot through its crawl") @click.option("--binary-id", help="Run one queued binary install directly on the bus") @click.option("--maintenance-only", is_flag=True, help="Only process sealed Snapshot maintenance and search-index backfills") @click.option("--no-stdin", is_flag=True, hidden=True, help="Run the scheduler even when stdin is not a TTY") def main( daemon: bool, crawl_id: str, snapshot_id: str, binary_id: str, maintenance_only: bool, no_stdin: bool, ): """ Process queued work. Modes: - No args + stdin piped: Process piped JSONL records - No args + TTY: Run the crawl runner for all work - --crawl-id: Run the crawl runner for that crawl only - --snapshot-id: Run one snapshot through its crawl only - --binary-id: Run one queued binary install directly on the bus """ from archivebox.core.shutdown_util import foreground_parent_watchdog, foreground_shutdown_signals if daemon and not snapshot_id and not binary_id and not crawl_id: try: os.environ[RUNNER_DAEMON_ENV] = "1" with ( foreground_shutdown_signals( on_signal=_exit_daemon_runner_on_signal, raise_on_first_signal=False, ), foreground_parent_watchdog(enabled=False), ): sys.exit(run_runner(daemon=True, maintenance_only=maintenance_only)) except KeyboardInterrupt: sys.exit(0) with foreground_shutdown_signals(), foreground_parent_watchdog(enabled=not daemon): if snapshot_id: sys.exit(run_snapshot_worker(snapshot_id)) if binary_id: try: from archivebox.services.runner import run_binary run_binary(binary_id) sys.exit(0) except KeyboardInterrupt: sys.exit(0) except Exception as e: rprint(f"[red]Runner error: {type(e).__name__}: {e}[/red]", file=sys.stderr) import traceback traceback.print_exc() sys.exit(1) if crawl_id: sys.exit( run_runner( daemon=False, crawl_id=crawl_id, maintenance_only=maintenance_only, ), ) if maintenance_only: sys.exit(run_runner(daemon=daemon, maintenance_only=True)) if not no_stdin and not sys.stdin.isatty(): sys.exit(process_stdin_records()) else: sys.exit(run_runner(daemon=daemon, maintenance_only=maintenance_only)) def run_snapshot_worker(snapshot_id: str) -> int: from archivebox.config import CONSTANTS from archivebox.core.takeover_util import enter_single_runner_gate from archivebox.core.shutdown_util import foreground_parent_watchdog, foreground_shutdown_signals from archivebox.machine.models import Process from archivebox.core.models import Snapshot from archivebox.services.runner import run_due_snapshot from django.utils import timezone current = Process.current() if not enter_single_runner_gate(current, data_dir=CONSTANTS.DATA_DIR): current.mark_exited() return 0 snapshot = None try: with foreground_shutdown_signals(), foreground_parent_watchdog(): for _ in range(10): snapshot = Snapshot.objects.select_related("crawl").get(id=snapshot_id) if snapshot.retry_at is None: snapshot.update_and_requeue(retry_at=timezone.now()) elif snapshot.retry_at > timezone.now(): break if not run_due_snapshot(snapshot, lock_seconds=60): break return 0 except KeyboardInterrupt: try: if snapshot is not None: snapshot.refresh_from_db() else: snapshot = Snapshot.objects.filter(id=snapshot_id).first() if snapshot is not None and snapshot.status != Snapshot.StatusChoices.SEALED: snapshot.update_and_requeue(retry_at=timezone.now()) except Exception: pass return 0 except Exception as e: rprint(f"[red]Runner error: {type(e).__name__}: {e}[/red]", file=sys.stderr) import traceback traceback.print_exc() return 1 finally: current.refresh_from_db() if current.status != Process.StatusChoices.EXITED: current.mark_exited() if __name__ == "__main__": main()