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.py | run_crawl()与CrawlRunner类 |
Crawl模型与队列转换 | archivebox/crawls/models.py | Crawl模型的原子队列状态迁移 |
Snapshot/ArchiveResult | archivebox/core/models.py | Snapshot、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 动作都只是"创建或选择数据库行",随后调用同一个 runner;runner 发出生命周期事件,abx-plugin 钩子完成实际抓取工作,service 投影器负责把Process与结果持久化。
源码印证:CrawlRunner 的初始化与事件总线装配
在 archivebox/services/runner.py 中,CrawlRunner.__init__完整展示了这条路径的装配顺序:
class CrawlRunner: def __init__(self, crawl, *, snapshot_ids=None, selected_plugins=None, ...): self.crawl = crawl self.bus = create_bus(name=_bus_name("ArchiveBox", str(crawl.id)), total_timeout=3600.0) self.catalog = get_plugin_catalog() HookProcessService(self.bus, emit_jsonl=False, interactive_tty=interactive_interrupts) register_sonic_daemon_event_handler(self.bus) PersistedProcessService(self.bus) ArchiveBoxBinaryService(self.bus) BinaryService(self.bus) TagService(self.bus) CrawlService(self.bus, crawl_id=str(crawl.id)) MachineService(self.bus) ... self.snapshot_service = SnapshotService(self.bus, crawl_id=str(crawl.id)) HookArchiveResultService(self.bus, emit_jsonl=False) ArchiveResultService(self.bus)可以看到:
- 每个
CrawlRunner都会创建一个独立的EventBus(create_bus,总超时 3600 秒),并注册多个投影服务:PersistedProcessService(持久化Process行)、ArchiveBoxBinaryService/BinaryService(二进制解析)、CrawlService(Crawl 事件投影)、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_install=True)→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 outputs"; Crawl.output_dir则落在CONSTANTS.USERS_DIR / created_by.username / CRAWLS_DIR_NAME / date / domain / crawl_id(archivebox/crawls/models.py)。
关于"旧版时间戳命名快照目录":Snapshot上带有fs_version字段(默认"0.9.0",core/models.py),用于区分不同文件系统版本;仓库中 tests/test_snapshot_filesystem_migration.py 专门覆盖了旧版目录向新布局迁移的场景。
Crawl 队列生命周期:显式状态机,无内存态漂移
Crawl的生命周期直接由 archivebox/crawls/models.py 中的Crawl类实现。文档特意强调了一个关键设计决策:
The database row is the durable state; the runner claims
retry_atwith a conditional update before it performs side effects, then calls the model's 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 的DefaultStatusChoices:queued/started/paused/sealed。Crawl在此基础上声明了各状态分组常量(crawls/models.py):
INITIAL_STATE = StatusChoices.QUEUED ACTIVE_STATE = StatusChoices.STARTED FINAL_STATES = (StatusChoices.SEALED,) RUNNABLE_STATES = (StatusChoices.QUEUED, StatusChoices.STARTED) INACTIVE_STATES = (StatusChoices.PAUSED, StatusChoices.SEALED)条件租约协议:claim → 副作用 → 显式迁移
ModelWithQueue(archivebox/workers/models.py)实现了统一的队列原语:
safe_update():带extra_filter的条件更新,更新行数不为 1 时记录SafeUpdateGuardMiss日志——这是防"陈旧写覆盖"的核心;claim_for_worker()/claim_processing_lock():UPDATE ... WHERE pk=? AND retry_at=? AND retry_at<=now(),把retry_at推进为"租约到期时间",返回是否抢到(CAS 语义);update_and_requeue():以extra_filter={"retry_at": self.retry_at}保证只有租约持有者能推进状态;pause()将retry_at置为RETRY_AT_MAX(datetime(9999,1,1)),使其退出可运行队列;resume()将其置回QUEUED并设置新的retry_at;ACTIVE_STATE_LEASE_SECONDS = 60,即活动租约默认 60 秒。
Crawl的具体生命周期方法(crawls/models.py):
mark_started():条件更新QUEUED → STARTED(extra_filter={"status": QUEUED}),并把retry_at设为now + 2s;seal():条件更新RUNNABLE_STATES → SEALED,置retry_at=None,随后schedule_child_snapshots_for_sealing()唤醒所有子快照、cleanup_runtime()清理 pid 文件与 persona 运行时;advance_lifecycle():runner 抢到行之后调用,按当前状态推进一步——QUEUED且已有快照全部完成则seal(),否则mark_started();STARTED且is_finished()则seal();pause()在暂停自身的同时会级联暂停所有非终态子快照;resume()会把所有PAUSED子快照批量置回QUEUED(crawls/models.py);cancel()先把 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_at=None): if self.kind == "update": run_scheduled_maintenance() # 直接运行维护,不创建 Crawl ... return None return self.enqueue(queued_at=queued_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_STATES(core/models.py),其中OPEN_STATES = (*RUNNABLE_STATES, PAUSED)是seal()允许的转换前置状态。
源码印证:start_processing 与 seal
Snapshot.start_processing()(core/models.py)实现QUEUED → STARTED的原子转换:
updated = ( type(self) .objects.filter(pk=self.pk, retry_at=owned_retry_at, status=self.StatusChoices.QUEUED) .update(status=self.StatusChoices.STARTED, retry_at=lease_until, modified_at=now) )seal()(core/models.py)则要求status__in=self.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中发现的子 URL(project_discovered_snapshots)、写入downloaded_at、处理crawl_max_size/crawl_timeout限额停止原因、条件更新SEALED、清空RETRY_PLUGINS并写出index.jsonl(snapshot.write_index_jsonl)。
这里的"投影"含义是:事件是 abx-dl 层的事实(facts),service 把这些事实单向写入 Django 模型行,而非模型主动轮询事件。
ArchiveResult 投影:事件驱动的结果行
ArchiveResult不由独立的内存状态机驱动。runner 创建 queued 行,ArchiveResultService把ArchiveResultEvent与ProcessCompletedEvent的数据投影(project)进这些行:
ArchiveResult.StatusChoices完整定义于 archivebox/core/models.py:
class 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 Results"FINAL_STATES = (SUCCEEDED, FAILED, SKIPPED, NORESULTS)——与图中四个终态一一对应。succeeded、failed、skipped、noresults均为终态;backoff表示"可恢复等待",恢复后回到started(图中recoverable wait → backoff → resumed work → started)。
源码印证:行结构唯一性约束与投影实现
每行记录产出它的插件与钩子,并存储结构化输出、文件元数据、耗时与错误详情(core/models.py):
snapshot(FK,on_delete=CASCADE)、plugin(插件名)、hook_name(如on_Snapshot__50_wget.py);process(OneToOneField → machine.Process,记录 cmd/pwd/stdout/stderr 等执行细节);- 输出字段:
output_str(人类可读摘要)、output_json(结构化元数据:headers、redirects 等)、output_files({相对路径: 元数据}字典)、output_size(总字节数)、output_mimetypes(按大小排序的 mimetype CSV); - 时间字段:
start_ts/end_ts;notes存放错误详情; - 唯一约束
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_id; get_or_create_by_hook幂等获取/创建结果行,diff 后只更新变化的字段;- 若结果到达
SUCCEEDED/NORESULTS,顺带更新快照标题(title插件的title.txt优先),并在urls.jsonl存在时调用project_discovered_snapshots()把解析器发现的新 URL 持久化为子快照。
测试印证
仓库测试覆盖了这套生命周期的核心路径,可作为进一步阅读入口:
- archivebox/tests/test_crawl_runner.py:runner 认领、执行与封存的完整流程;
- archivebox/tests/test_crawl_service.py:Crawl 事件投影;
- archivebox/tests/test_snapshot_service.py:Snapshot 事件投影与
finalize_completed_snapshot; - archivebox/tests/test_process_service.py:Process 行持久化与状态;
- archivebox/tests/test_resource_admission.py:runner 启动时的资源准入(内存/磁盘/网络限制下的 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/bin,LIB_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),仅供参考