From 04f20447f868c4e3f035841ef7cd672869bca540 Mon Sep 17 00:00:00 2001 From: Nick Sweeting Date: Thu, 28 May 2026 20:32:51 -0700 Subject: [PATCH] Fix snapshot finalization race --- archivebox/core/models.py | 8 +-- archivebox/crawls/models.py | 7 +- archivebox/services/runner.py | 7 +- archivebox/services/snapshot_service.py | 85 +++++++++++++------------ 4 files changed, 56 insertions(+), 51 deletions(-) diff --git a/archivebox/core/models.py b/archivebox/core/models.py index d4f45097..3dc82c14 100755 --- a/archivebox/core/models.py +++ b/archivebox/core/models.py @@ -3199,11 +3199,9 @@ class SnapshotMachine(BaseStateMachine): # Clean up background hooks self.snapshot.cleanup() - self.snapshot.update_and_requeue( - retry_at=None, - status=Snapshot.StatusChoices.SEALED, - refresh=False, - ) + self.snapshot.status = Snapshot.StatusChoices.SEALED + self.snapshot.retry_at = None + self.snapshot.save(update_fields=["status", "retry_at", "modified_at"]) # Crawl finalization is handled by the runner/CrawlService cleanup # phase. Sealing the parent crawl here races recursive discovery: diff --git a/archivebox/crawls/models.py b/archivebox/crawls/models.py index ae3fb1df..b104cdf6 100755 --- a/archivebox/crawls/models.py +++ b/archivebox/crawls/models.py @@ -1540,10 +1540,9 @@ class CrawlMachine(BaseStateMachine): # Clean up background hooks and run on_CrawlEnd hooks self.crawl.cleanup() - self.crawl.update_and_requeue( - retry_at=None, - status=Crawl.StatusChoices.SEALED, - ) + self.crawl.status = Crawl.StatusChoices.SEALED + self.crawl.retry_at = None + self.crawl.save(update_fields=["status", "retry_at", "modified_at"]) # ============================================================================= diff --git a/archivebox/services/runner.py b/archivebox/services/runner.py index 44fc3aea..1d60dd43 100644 --- a/archivebox/services/runner.py +++ b/archivebox/services/runner.py @@ -68,7 +68,7 @@ from .binary_service import BinaryService from .crawl_service import CrawlService from .machine_service import MachineService from .process_service import ProcessService as PersistedProcessService -from .snapshot_service import SnapshotService +from .snapshot_service import SnapshotService, finalize_completed_snapshot from .tag_service import TagService @@ -1045,6 +1045,11 @@ class CrawlRunner: raise RuntimeError(f"Snapshot {snapshot_id} did not complete") await completed_snapshot.wait(timeout=snapshot_phase_timeout) await completed_snapshot.event_results_list() + # SnapshotCompletedEvent is the normal projection path, but the + # runner is the scheduler owner. Finalize idempotently here too + # so a completed snapshot cannot remain STARTED if the event was + # observed before its DB projector advanced the state machine. + await sync_to_async(finalize_completed_snapshot, thread_sensitive=True)(snapshot_id) if snapshot["status"] == "sealed": await sync_to_async(run_snapshot_maintenance, thread_sensitive=True)(snapshot_id) return diff --git a/archivebox/services/snapshot_service.py b/archivebox/services/snapshot_service.py index a6543212..a983c1b3 100644 --- a/archivebox/services/snapshot_service.py +++ b/archivebox/services/snapshot_service.py @@ -3,14 +3,56 @@ from __future__ import annotations import sys from asgiref.sync import sync_to_async -from django.core.exceptions import ValidationError from django.utils import timezone +from django.core.exceptions import ValidationError from rich import print as rprint from abx_dl.events import SnapshotCompletedEvent, SnapshotEvent from abx_dl.limits import CrawlLimitState from abx_dl.services.base import BaseService +def finalize_completed_snapshot(snapshot_id: str) -> None: + from archivebox.core.models import Snapshot + + snapshot = Snapshot.objects.select_related("crawl", "crawl__created_by").filter(id=snapshot_id).first() + if snapshot is None: + return + + if snapshot.downloaded_at is None: + snapshot.downloaded_at = timezone.now() + snapshot.save(update_fields=["downloaded_at", "modified_at"]) + + stop_reason = _crawl_limit_stop_reason(snapshot.crawl) + if snapshot.crawl_id and stop_reason in ("crawl_max_size", "crawl_timeout"): + Snapshot.objects.filter( + crawl_id=snapshot.crawl_id, + status=Snapshot.StatusChoices.QUEUED, + ).exclude(id=snapshot.id).update( + status=Snapshot.StatusChoices.SEALED, + retry_at=None, + modified_at=timezone.now(), + ) + + if snapshot.status == Snapshot.StatusChoices.QUEUED: + snapshot.sm.tick() + snapshot.refresh_from_db() + if snapshot.status == Snapshot.StatusChoices.STARTED and snapshot.is_finished_processing(): + snapshot.sm.seal() + snapshot.refresh_from_db() + + snapshot.write_index_jsonl() + snapshot.write_json_details() + snapshot.write_html_details() + + +def _crawl_limit_stop_reason(crawl) -> str: + from archivebox.config.common import get_config + + config = get_config(crawl=crawl, include_machine=False) + config["CRAWL_DIR"] = str(crawl.output_dir) + return CrawlLimitState.from_config(config).get_stop_reason() + + class SnapshotService(BaseService): LISTENS_TO = [SnapshotEvent, SnapshotCompletedEvent] EMITS = [] @@ -54,43 +96,4 @@ class SnapshotService(BaseService): await sync_to_async(snapshot.ensure_crawl_symlink, thread_sensitive=True)() async def on_SnapshotCompletedEvent(self, event: SnapshotCompletedEvent) -> None: - from archivebox.core.models import Snapshot - - snapshot = await Snapshot.objects.select_related("crawl", "crawl__created_by").filter(id=event.snapshot_id).afirst() - snapshot_id: str | None = None - if snapshot is not None: - snapshot.downloaded_at = snapshot.downloaded_at or timezone.now() - await snapshot.asave(update_fields=["downloaded_at", "modified_at"]) - stop_reason = await sync_to_async(self._crawl_limit_stop_reason, thread_sensitive=True)(snapshot.crawl) - if snapshot.crawl_id and stop_reason in ("crawl_max_size", "crawl_timeout"): - await ( - Snapshot.objects.filter( - crawl_id=snapshot.crawl_id, - status=Snapshot.StatusChoices.QUEUED, - ) - .exclude(id=snapshot.id) - .aupdate( - status=Snapshot.StatusChoices.SEALED, - retry_at=None, - modified_at=timezone.now(), - ) - ) - if snapshot.status == Snapshot.StatusChoices.QUEUED: - await sync_to_async(snapshot.sm.tick, thread_sensitive=True)() - await sync_to_async(snapshot.refresh_from_db, thread_sensitive=True)() - if snapshot.status == Snapshot.StatusChoices.STARTED: - await sync_to_async(snapshot.sm.seal, thread_sensitive=True)() - snapshot_id = str(snapshot.id) - if snapshot_id: - snapshot = await Snapshot.objects.filter(id=snapshot_id).select_related("crawl", "crawl__created_by").afirst() - if snapshot is not None: - await sync_to_async(snapshot.write_index_jsonl, thread_sensitive=True)() - await sync_to_async(snapshot.write_json_details, thread_sensitive=True)() - await sync_to_async(snapshot.write_html_details, thread_sensitive=True)() - - def _crawl_limit_stop_reason(self, crawl) -> str: - from archivebox.config.common import get_config - - config = get_config(crawl=crawl, include_machine=False) - config["CRAWL_DIR"] = str(crawl.output_dir) - return CrawlLimitState.from_config(config).get_stop_reason() + await sync_to_async(finalize_completed_snapshot, thread_sensitive=True)(event.snapshot_id)