ArchiveBox/archivebox/cli/archivebox_run.py
2026-09-02 12:56:23 -07:00

455 lines
17 KiB
Python

#!/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()