
ArchiveBox 架构图解执行链路、持久化数据与 Crawl/Snapshot 队列生命周期【免费下载链接】ArchiveBox Open source self-hosted web archiving. Takes URLs/browser history/bookmarks/Pocket/Pinboard/etc., saves HTML, JS, PDFs, media, and more...项目地址: https://gitcode.com/gh_mirrors/ar/ArchiveBox本文以 docs/ArchiveBox-Architecture-Diagrams.md 为骨架结合 ArchiveBox 仓库内的archivebox/services/runner.py、archivebox/crawls/models.py、archivebox/core/models.py、archivebox/workers/models.py等源码实现完整讲解 ArchiveBox 当前唯一的正常抓取执行路径、持久化数据布局以及Crawl、Snapshot、ArchiveResult三个核心模型的队列状态机与生命周期转换。读完本文你将掌握一次抓取Crawl从入口到落盘的完整调用链、数据库行如何作为唯一可信状态source of truth、runner 如何通过条件更新conditional update抢占retry_at租约、abxpkg 如何完成二进制解析与安装以及各状态迁移在源码中的具体落点。模块地图架构实现主要分布在哪里文档开头给出了一张执行与持久化路径的地图指明 ArchiveBox 当前架构的核心实现位置关注点主要实现位置职责CLI 入口archivebox/cli/archivebox add、archivebox update、archivebox schedule等命令行入口抓取与快照执行archivebox/services/runner.pyrun_crawl()与CrawlRunner类Crawl模型与队列转换archivebox/crawls/models.pyCrawl模型的原子队列状态迁移Snapshot/ArchiveResultarchivebox/core/models.pySnapshot、ArchiveResult模型及快照队列转换总线事件投影器archivebox/services/CrawlService、SnapshotService、ArchiveResultService、ProcessService等 bus event projectors二进制解析与插件钩子abxpkg与abx-plugins二进制解析abxpkg与插件钩子系统从源码结构看这一分工遵循模型负责数据库状态转换、runner 负责执行编排、service 负责把总线事件投影为持久化行的职责分离Crawl/Snapshot通过ModelWithQueue继承统一的队列字段与租约协议runner 不直接写状态而是调用模型暴露的显式生命周期方法。高层执行流唯一的一条抓取路径文档给出了 ArchiveBox 的高层执行流程图文档明确强调ArchiveBox 只有一条正常的抓取执行路径。CLI 命令与 Web/API 动作都只是创建或选择数据库行随后调用同一个 runnerrunner 发出生命周期事件abx-plugin 钩子完成实际抓取工作service 投影器负责把Process与结果持久化。源码印证CrawlRunner 的初始化与事件总线装配在 archivebox/services/runner.py 中CrawlRunner.__init__完整展示了这条路径的装配顺序class CrawlRunner: def __init__(self, crawl, *, snapshot_idsNone, selected_pluginsNone, ...): self.crawl crawl self.bus create_bus(name_bus_name(ArchiveBox, str(crawl.id)), total_timeout3600.0) self.catalog get_plugin_catalog() HookProcessService(self.bus, emit_jsonlFalse, interactive_ttyinteractive_interrupts) register_sonic_daemon_event_handler(self.bus) PersistedProcessService(self.bus) ArchiveBoxBinaryService(self.bus) BinaryService(self.bus) TagService(self.bus) CrawlService(self.bus, crawl_idstr(crawl.id)) MachineService(self.bus) ... self.snapshot_service SnapshotService(self.bus, crawl_idstr(crawl.id)) HookArchiveResultService(self.bus, emit_jsonlFalse) ArchiveResultService(self.bus)可以看到每个CrawlRunner都会创建一个独立的EventBuscreate_bus总超时 3600 秒并注册多个投影服务PersistedProcessService持久化Process行、ArchiveBoxBinaryService/BinaryService二进制解析、CrawlServiceCrawl 事件投影、SnapshotService快照生命周期、ArchiveResultService结果投影、TagService标签同步以及 Sonic 搜索守护进程的事件处理器。runner 通过enqueue_snapshot()/run_snapshot()/wait_for_snapshot_tasks()等方法以 asyncio 任务方式调度快照并通过snapshot_semaphore与CRAWL_MAX_CONCURRENT_SNAPSHOTS控制并发runner.py。run()在启动时先做资源准入检查defer_crawl_for_resources见 archivebox/services/resource_admission.py然后load_run_state()加载或创建初始快照进入run_crawl()主流程。源码印证run_crawl 中的钩子编排run_crawl()runner.py是执行流的中枢加载快照 payload 并规范化运行时配置normalize_runtime_config设置ABX_RUNTIME archivebox按相计算超时compute_phase_timeout(CrawlSetup hooks)、Snapshot阶段超时 120s、CrawlCleanup阶段超时最终合成 crawl 生命周期总超时依次发出MachineEvent用户配置与派生配置→InstallEvent安装阶段通过PluginBinariesService处理auto_installTrue→CrawlSetupEvent→CrawlStartEvent→CrawlCleanupEvent→CrawlCompletedEvent注册两个内部事件处理器on_archivebox_CrawlStartEvent__run_snapshots按并发槽位批量 enqueue 快照与on_archivebox_CrawlEvent__run_recursive_crawl执行 Setup/Start/Cleanup 全阶段并监听取消watch_for_cancelled_crawl每 1 秒轮询一次数据库若 crawl 已被置为 SEALED则发出CrawlAbortEvent中止。这与文档runner 发出生命周期事件abx-plugin 钩子完成抓取工作service 投影器持久化进程与结果的描述完全吻合。二进制解析abxpkg 统一入口文档明确指出二进制发现与安装一律经由 abxpkg。优先使用宿主机上兼容的二进制Compatible host binary托管安装Managed install fallback作为兜底解析成功的二进制会被投影到LIB_DIR/env/bin供程序化调用而LIB_DIR/bin仅是人类使用的便利目录。在 archivebox/config/constants.py 中DEFAULT_ABXPKG_LIB_DIR默认为user_config_path(abx) / lib可通过ABXPKG_LIB_DIR环境变量覆盖。runner 运行结束后project_abxpkg_derived_cache_to_db()会把 abxpkg 解析出的二进制缓存回数据库runner.py供archivebox系列 CLI 与machine相关命令查询。持久化数据布局数据库为唯一事实源文档给出了数据目录的结构图数据库是模型状态的唯一事实源快照目录存放抓取产物captured artifacts与渲染后的元数据旧版集合中可能还存在以时间戳命名的历史快照目录legacy timestamp-named snapshot directories。源码印证目录常量定义在 archivebox/config/constants.py 中ARCHIVE_DIR_NAME: str archive USERS_DIR_NAME: str users SNAPSHOTS_DIR_NAME: str snapshots CRAWLS_DIR_NAME: str crawls SOURCES_DIR_NAME: str sources LOGS_DIR_NAME: str logs ... SOURCES_DIR: Path DATA_DIR / SOURCES_DIR_NAME LOGS_DIR: Path DATA_DIR / LOGS_DIR_NAME ... SQL_INDEX_FILENAME: str index.sqlite3 DATABASE_FILE: Path DATA_DIR / SQL_INDEX_FILENAMEindex.sqlite3即DATABASE_FILE DATA_DIR / index.sqlite3是所有模型行的落点快照目录路径由Snapshot.output_dir计算见 archivebox/core/models.py 的fs_version与get_storage_path_for_version()逻辑其形态即为archive/users/user/snapshots/date/domain/uuid/目录内是插件命名空间子目录如wget/、singlefile/、pdf/、title/对应图中的 Plugin-namespaced outputsCrawl.output_dir则落在CONSTANTS.USERS_DIR / created_by.username / CRAWLS_DIR_NAME / date / domain / crawl_idarchivebox/crawls/models.py。关于旧版时间戳命名快照目录Snapshot上带有fs_version字段默认0.9.0core/models.py用于区分不同文件系统版本仓库中 tests/test_snapshot_filesystem_migration.py 专门覆盖了旧版目录向新布局迁移的场景。Crawl 队列生命周期显式状态机无内存态漂移Crawl的生命周期直接由 archivebox/crawls/models.py 中的Crawl类实现。文档特意强调了一个关键设计决策The database row is the durable state; the runner claimsretry_atwith a conditional update before it performs side effects, then calls the models explicit lifecycle methods. There is deliberately no second in-memory state machine that can drift from the row owned by another process.即数据库行是持久状态runner 在产生副作用前先用条件更新抢占retry_at再调用模型显式的生命周期方法刻意不引入第二套可能与其他进程所持行发生漂移的内存状态机。状态集合来自 archivebox/workers/models.py 的DefaultStatusChoicesqueued/started/paused/sealed。Crawl在此基础上声明了各状态分组常量crawls/models.pyINITIAL_STATE StatusChoices.QUEUED ACTIVE_STATE StatusChoices.STARTED FINAL_STATES (StatusChoices.SEALED,) RUNNABLE_STATES (StatusChoices.QUEUED, StatusChoices.STARTED) INACTIVE_STATES (StatusChoices.PAUSED, StatusChoices.SEALED)条件租约协议claim → 副作用 → 显式迁移ModelWithQueuearchivebox/workers/models.py实现了统一的队列原语safe_update()带extra_filter的条件更新更新行数不为 1 时记录SafeUpdateGuardMiss日志——这是防陈旧写覆盖的核心claim_for_worker()/claim_processing_lock()UPDATE ... WHERE pk? AND retry_at? AND retry_atnow()把retry_at推进为租约到期时间返回是否抢到CAS 语义update_and_requeue()以extra_filter{retry_at: self.retry_at}保证只有租约持有者能推进状态pause()将retry_at置为RETRY_AT_MAXdatetime(9999,1,1)使其退出可运行队列resume()将其置回QUEUED并设置新的retry_atACTIVE_STATE_LEASE_SECONDS 60即活动租约默认 60 秒。Crawl的具体生命周期方法crawls/models.pymark_started()条件更新QUEUED → STARTEDextra_filter{status: QUEUED}并把retry_at设为now 2sseal()条件更新RUNNABLE_STATES → SEALED置retry_atNone随后schedule_child_snapshots_for_sealing()唤醒所有子快照、cleanup_runtime()清理 pid 文件与 persona 运行时advance_lifecycle()runner 抢到行之后调用按当前状态推进一步——QUEUED且已有快照全部完成则seal()否则mark_started()STARTED且is_finished()则seal()pause()在暂停自身的同时会级联暂停所有非终态子快照resume()会把所有PAUSED子快照批量置回QUEUEDcrawls/models.pycancel()先把 crawl 置为SEALED并唤醒子快照再由 runner 认领SEALED到期的行执行清理钩子——文档所述explicit seal路径。关于暂停时也调度子快照暂停、恢复时回到可运行队列从源码看Crawl.pause()调用schedule_child_snapshots_for_pause()crawls/models.py它只把子快照的retry_at唤醒为now真正的PAUSED转换由每个Snapshot的 runner 认领后在reconcile_parent_lifecycle()core/models.py中完成——父级只叫醒子行不做大事务这就是文档强调的runner claim performs the real pause transition。调度维护CrawlSchedule 直接分发文档补充说明定时维护由CrawlSchedule直接分发不会创建合成synthetic的 crawl 或 snapshot。在 archivebox/crawls/models.py 中def dispatch(self, queued_atNone): if self.kind update: run_scheduled_maintenance() # 直接运行维护不创建 Crawl ... return None return self.enqueue(queued_atqueued_at) # 否则入队一条普通 CrawlCrawlSchedule通过is_due()is_enabled and next_run_at now判断是否到期enqueue()从模板行复制urls/max_depth/persona等字段并生成config配置快照新建一条状态为QUEUED的普通 Crawl——与文档描述完全一致。Snapshot 队列生命周期与 Crawl 相同的租约协议Snapshot的生命周期由 archivebox/core/models.py 中的Snapshot类实现使用与Crawl相同的条件retry_at认领协议Snapshot同样声明RUNNABLE_STATES/OPEN_STATEScore/models.py其中OPEN_STATES (*RUNNABLE_STATES, PAUSED)是seal()允许的转换前置状态。源码印证start_processing 与 sealSnapshot.start_processing()core/models.py实现QUEUED → STARTED的原子转换updated ( type(self) .objects.filter(pkself.pk, retry_atowned_retry_at, statusself.StatusChoices.QUEUED) .update(statusself.StatusChoices.STARTED, retry_atlease_until, modified_atnow) )seal()core/models.py则要求status__inself.OPEN_STATES且retry_at仍为自己持有的租约值成功后调用finalize_output_metadata()汇总输出元数据。advance_lifecycle()core/models.py刻意保持简单QUEUED → STARTED要求 URL 就绪其余状态不推进——快照是否完成由 abx-dl 发出的SnapshotCompletedEvent决定ArchiveResult的投影状态绝不反过来驱动 Snapshot 完成。这与文档runner 为每个选中钩子创建一条 queued 的 ArchiveResult……在所有结果到达终态后封存快照的描述呼应同时也说明了窄范围的搜索索引维护操作对已封存快照执行是刻意保留的例外不会重新打开或发明第二条通用生命周期路径——对应 runner.py 中allow_maintenance_on_inactive_crawl的显式限定initial_snapshot_ids selected_plugins crawl 已 SEALED。SnapshotService事件投影与租约续期archivebox/services/snapshot_service.py 中的SnapshotService监听SnapshotEvent与SnapshotCompletedEvent收到SnapshotEvent时调用snapshot.advance_lifecycle()完成QUEUED → STARTED并把 (event_id, retry_at, was_sealed, retry_plugins) 记录在 runner 进程内的_run_ownership字典中renew_lease()每 10 秒由 runner 心跳调用runner.py以条件更新把retry_at续到lease_until now ACTIVE_STATE_LEASE_SECONDS续约失败说明租约被他人接管runner 会取消对应快照任务finalize_completed_snapshot()在SnapshotCompletedEvent到达后执行投影urls.jsonl中发现的子 URLproject_discovered_snapshots、写入downloaded_at、处理crawl_max_size/crawl_timeout限额停止原因、条件更新SEALED、清空RETRY_PLUGINS并写出index.jsonlsnapshot.write_index_jsonl。这里的投影含义是事件是 abx-dl 层的事实factsservice 把这些事实单向写入 Django 模型行而非模型主动轮询事件。ArchiveResult 投影事件驱动的结果行ArchiveResult不由独立的内存状态机驱动。runner 创建 queued 行ArchiveResultService把ArchiveResultEvent与ProcessCompletedEvent的数据投影project进这些行ArchiveResult.StatusChoices完整定义于 archivebox/core/models.pyclass StatusChoices(models.TextChoices): QUEUED queued, Queued STARTED started, Started PAUSED paused, Paused BACKOFF backoff, Waiting to retry SUCCEEDED succeeded, Succeeded FAILED failed, Failed SKIPPED skipped, Skipped NORESULTS noresults, No ResultsFINAL_STATES (SUCCEEDED, FAILED, SKIPPED, NORESULTS)——与图中四个终态一一对应。succeeded、failed、skipped、noresults均为终态backoff表示可恢复等待恢复后回到started图中recoverable wait → backoff → resumed work → started。源码印证行结构唯一性约束与投影实现每行记录产出它的插件与钩子并存储结构化输出、文件元数据、耗时与错误详情core/models.pysnapshotFKon_deleteCASCADE、plugin插件名、hook_name如on_Snapshot__50_wget.pyprocessOneToOneField → machine.Process记录 cmd/pwd/stdout/stderr 等执行细节输出字段output_str人类可读摘要、output_json结构化元数据headers、redirects 等、output_files{相对路径: 元数据}字典、output_size总字节数、output_mimetypes按大小排序的 mimetype CSV时间字段start_ts/end_tsnotes存放错误详情唯一约束UniqueConstraint(fields[snapshot, plugin, hook_name])保证每个快照的每个钩子恰好一行结果配合get_or_create_by_hook()core/models.py实现幂等投影。投影逻辑在 archivebox/services/archive_result_service.py 的_save_archiveresult_event_to_db()中按event.snapshot_id查快照select_related(crawl, crawl__created_by)解析输出元数据优先用事件中的OutputManifest否则扫描插件输出目录OutputManifest.scan(plugin_dir, ...)通过ProcessStartedEvent反查Process行按pwd cmd started_at可选pid并关联process_idget_or_create_by_hook幂等获取/创建结果行diff 后只更新变化的字段若结果到达SUCCEEDED/NORESULTS顺带更新快照标题title插件的title.txt优先并在urls.jsonl存在时调用project_discovered_snapshots()把解析器发现的新 URL 持久化为子快照。测试印证仓库测试覆盖了这套生命周期的核心路径可作为进一步阅读入口archivebox/tests/test_crawl_runner.pyrunner 认领、执行与封存的完整流程archivebox/tests/test_crawl_service.pyCrawl 事件投影archivebox/tests/test_snapshot_service.pySnapshot 事件投影与finalize_completed_snapshotarchivebox/tests/test_process_service.pyProcess 行持久化与状态archivebox/tests/test_resource_admission.pyrunner 启动时的资源准入内存/磁盘/网络限制下的 defer 逻辑。设计要点小结把文档与源码对照后可以提炼出这套架构的几个核心设计原则数据库行即状态杜绝双状态机。Crawl/Snapshot的状态转换全部经由ModelWithQueue的条件更新原语safe_update/claim_for_worker/update_and_requeue完成runner 进程内的_run_ownership等内存结构只记录本次运行认领了哪一行不构成可漂移的第二套状态机。先认领、后副作用。runner 先通过retry_at条件更新抢到租约才执行抓取副作用心跳续约默认 60 秒租约、10 秒一次续约保证长时间任务不被误判为超时同时允许未来 PostgreSQL 多机 runner 以这些认领为协调边界源码注释中已明确此意图。事件驱动投影。abx-dl 层发出生命周期事件CrawlEvent、SnapshotEvent、ArchiveResultEvent、ProcessEvent等archivebox/services/下的 service 类单向投影为 Django 行ArchiveResult投影状态不会反过来决定 Snapshot 是否封存。二进制解析统一走 abxpkg。宿主机兼容二进制优先托管安装兜底解析结果投影到LIB_DIR/env/binLIB_DIR/bin仅供人类便利使用。如果你希望进一步验证这些结论可以依次阅读 docs/ArchiveBox-Architecture-Diagrams.md 原始文档、archivebox/services/runner.py、archivebox/crawls/models.py、archivebox/core/models.py 与 archivebox/workers/models.py并结合archivebox/tests/下的 runner/service 测试用例按图索骥。【免费下载链接】ArchiveBox Open source self-hosted web archiving. Takes URLs/browser history/bookmarks/Pocket/Pinboard/etc., saves HTML, JS, PDFs, media, and more...项目地址: https://gitcode.com/gh_mirrors/ar/ArchiveBox创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考