Files
ComfyUI/app/assets/scanner_changes.py
Simon Pinfold f427c3a285 Move file reads out of the remaining asset write transactions (#16486)
* Take the SQLite write lock up front for scan and output-registration writes

The scanner's seeding and reference sync, and executed-output registration,
write through a separate engine whose transactions open with BEGIN
IMMEDIATE, and the database runs in WAL mode. On those paths stat, hashing
and metadata extraction now happen before the write transaction opens, and
reference-sync results are applied only to rows unchanged since they were
observed. busy_timeout stays at pysqlite's 5s default.

A fast-scan batch now commits as one transaction, so an unexpected error
partway through discards the whole batch; the next scan recreates it.
Enrichment, verification, uploads and tagging still write through the
existing sessions.

Migration backups use SQLite's backup API, and a legacy database is
checkpointed before it is relocated, since in WAL mode committed pages can
live in the -wal file that a plain file copy misses.

* Skip relocating a legacy database whose WAL cannot be checkpointed

The checkpoint can report busy without raising; moving the file then would
leave committed pages behind in the -wal. Also keep the source's file mode
on SQLite backups, as the plain copy did.

* Replace run_write_txn with a create_write_session factory

Write paths open the write engine's session the same way the rest of the
code opens create_session(), and commit explicitly. Executed-output
registration reads the new record's fields before committing, so expiry
does not reload them in a second write transaction.

* Document the scanner's pre-transaction observation types

* Document that write sessions must not nest

* Warn that a nested write session looks like lock contention

* Seed a hashed spec with the stat its hash was verified against

* Bound the SQLite backup so a locked destination cannot hang startup

* Time out a backup only while it is blocked

* Commit each drained entry before hashing the next

drain_pending_verifications and drain_transition_queue ran every entry in
one session, so an entry that wrote (marking a vanished file missing, say)
left a transaction open while the next entry's file was stat'ed and
hashed. Commit at the top of each entry instead, so the hash runs with no
transaction open.

* Seed settled watch-list entries through insert_asset_specs

tick_watch_list seeded each settled file through seed_asset_specs in the
caller's session, where the enrich phase still held drain_pending's
writes open, and each seed ran in a deferred savepoint that reads before
it writes. Collect the settled specs and hand them to insert_asset_specs,
which stats and hashes before opening one write session. The seeder
commits the pending verifications before ticking, so no transaction is
open while the watch list stats or waits for the write lock.

* Read upload metadata before the claim transaction

_create_upload_record read the file for system metadata after the
content claim had opened a write transaction. Callers now extract the
metadata before opening their session (or, when reusing content, before
claiming it) and pass it in. The claim's own stat re-check stays inside
the transaction: that is what makes the claim sound.

* Keep the upgrade error when restoring the backup also fails

If restoring the pre-upgrade backup raised, that exception replaced the
upgrade's, and the backup's location was never logged. Log the upgrade
error first, log where the pre-upgrade copy is kept if the restore or its
cleanup fails, and re-raise the upgrade error either way.

* Note the drains' session requirement and word the restore log for either failure

The drains' per-entry commit only leaves no transaction open on a
create_session() session. The restore log now covers a failed backup
removal as well as a failed restore. The watch-list admission test
patches insert_asset_specs, the seam tick_watch_list now calls.

* Log the real error when a seed spec cannot be read

observe_asset_specs treated every OSError as a vanished file, so a
permission error or an I/O error on a file that still exists was
reported only as "Skipping vanished asset during scan". A missing file
is still handled as before; any other OSError is now also logged through
_log_scan_error before the spec is skipped. The vanished-path test's
fake now raises FileNotFoundError, the error a vanished file produces.

* Report an unreadable seed spec once, without also calling it vanished

* Skip a pending verification whose row changed while it was hashed

Committing per entry means no lock is held while the file is hashed, so
another writer can retire or replace the row in that time. The drain
would then store a hash on a missing row, or split it and attach a
record to the other writer's row. Re-read the row after hashing and skip
it unless it is still live with the hash, size and mtime it was loaded
with, as apply_reference_observations does.

* Remove the stored upload file in the reupload claim test

The test cleaned up its temp files but left the first upload's stored
file in the output directory.
2026-09-23 17:23:56 -07:00

218 lines
7.7 KiB
Python

"""Reconciles catalogued content against what is actually on disk: retiring rows
whose file is gone, splitting a row whose bytes changed, and recovering one
whose file came back. Recovery fires only when the returning file's hash
identifies exactly one missing row and no live row already occupies that path,
so a restored file can never leave two live rows describing one location.
"""
from __future__ import annotations
import os
from pathlib import Path
from typing import Literal
import sqlalchemy as sa
from sqlalchemy.orm import Session
from app.assets.database.models import AssetContent
from app.assets.database.queries.records import (
create_content,
create_record,
mark_content_missing,
unset_content_missing,
)
from app.assets.helpers import sql_path_under_prefix, to_stored_hash
from app.assets.services.path_utils import compute_loader_path, get_name_and_tags_from_asset_path
from app.assets.services.snapshot_hash import snapshot_hash
_pending_verification_ids: list[str] = []
_pending_recovery_paths: list[str] = []
def clear_pending_verifications() -> None:
_pending_verification_ids.clear()
_pending_recovery_paths.clear()
def queue_pending_verification(content_id: str) -> None:
if content_id not in _pending_verification_ids:
_pending_verification_ids.append(content_id)
def pending_recovery_count() -> int:
return len(_pending_recovery_paths)
def recover_missing_content(
session: Session,
path: str,
snapshot: tuple[str, os.stat_result] | None,
hashing_is_enabled: bool,
) -> Literal["recovered", "no_match", "unstable"]:
"""``snapshot`` is ``snapshot_hash(path)``, taken by the caller before its write
transaction opens so the file is never hashed while the write lock is held."""
if not hashing_is_enabled:
return "no_match"
occupied = session.scalar(
sa.select(AssetContent.id)
.where(AssetContent.path == path, AssetContent.is_missing.is_(False))
.limit(1)
)
if occupied is not None:
return "no_match"
if snapshot is None:
if path not in _pending_recovery_paths:
_pending_recovery_paths.append(path)
return "unstable"
digest, verified_stat = snapshot
stored_hash = to_stored_hash(digest)
matches = list(
session.scalars(
sa.select(AssetContent).where(
AssetContent.path == path,
AssetContent.is_missing.is_(True),
AssetContent.hash == stored_hash,
)
)
)
if len(matches) == 1:
recovered = matches[0]
unset_content_missing(session, recovered.id)
recovered.size_bytes = verified_stat.st_size
recovered.mtime_ns = verified_stat.st_mtime_ns
return "recovered"
if len(matches) > 1:
return "no_match"
null_hash_matches = list(
session.scalars(
sa.select(AssetContent).where(
AssetContent.path == path,
AssetContent.is_missing.is_(True),
AssetContent.hash.is_(None),
)
)
)
if len(null_hash_matches) != 1:
return "no_match"
candidate = null_hash_matches[0]
if (candidate.size_bytes, candidate.mtime_ns) != (
verified_stat.st_size,
verified_stat.st_mtime_ns,
):
return "no_match"
unset_content_missing(session, candidate.id)
candidate.hash = stored_hash
candidate.size_bytes = verified_stat.st_size
candidate.mtime_ns = verified_stat.st_mtime_ns
return "recovered"
def is_path_under_prefixes(path: str, prefixes: list[str]) -> bool:
candidate = Path(os.path.abspath(path))
return any(candidate.is_relative_to(os.path.abspath(prefix)) for prefix in prefixes)
def split_content(session: Session, content: AssetContent, stat_result: os.stat_result, hash_value: str | None) -> AssetContent:
mark_content_missing(session, content.id)
name, tags = get_name_and_tags_from_asset_path(content.path)
replacement = create_content(
session,
path=content.path,
hash=hash_value,
size_bytes=stat_result.st_size,
mtime_ns=stat_result.st_mtime_ns,
)
create_record(
session,
content_id=replacement.id,
name=name,
loader_path=compute_loader_path(content.path),
tags=tags,
)
return replacement
def detect_content_change(
session: Session,
content: AssetContent,
stat_result: os.stat_result,
hashing_is_enabled: bool,
) -> None:
if content.mtime_ns == stat_result.st_mtime_ns:
# Ruling #10: size drift with unchanged mtime is undefined behavior.
return
if hashing_is_enabled:
queue_pending_verification(content.id)
return
if content.size_bytes == stat_result.st_size:
# User identity rule: a same-size mtime bump (rsync, cloud sync, backup restore) is the
# same file — never split, or the record's tags and metadata are destroyed.
# The stored hash goes with the refreshed stat: OFF mode cannot prove the bytes, and a
# refreshed stat alone would re-qualify the row to be served under a digest it may no
# longer match.
content.size_bytes = stat_result.st_size
content.mtime_ns = stat_result.st_mtime_ns
content.hash = None
return
split_content(session, content, stat_result, hash_value=None)
def drain_pending_verifications(session: Session, limit: int | None = None) -> int:
queued_count = min(len(_pending_verification_ids), limit or len(_pending_verification_ids))
processed = 0
for _ in range(queued_count):
# Commit the previous entry's writes so this entry's hash runs with no transaction open.
# That needs a create_session() session: on a write session the next read takes the lock.
session.commit()
content_id = _pending_verification_ids.pop(0)
content = session.get(AssetContent, content_id)
if content is None or content.is_missing:
continue
loaded = (content.hash, content.size_bytes, content.mtime_ns)
try:
os.stat(content.path, follow_symlinks=True)
except FileNotFoundError:
mark_content_missing(session, content.id)
processed += 1
continue
except OSError:
queue_pending_verification(content_id)
continue
try:
snapshot = snapshot_hash(content.path)
except OSError:
queue_pending_verification(content_id)
continue
if snapshot is None:
queue_pending_verification(content_id)
continue
digest, verified_stat = snapshot
stored_hash = to_stored_hash(digest)
# Skip a row another writer retired or changed while the file was hashed.
content = session.get(AssetContent, content_id, populate_existing=True)
if content is None or content.is_missing or (content.hash, content.size_bytes, content.mtime_ns) != loaded:
continue
if content.hash == stored_hash or content.hash is None:
content.hash = stored_hash
content.size_bytes = verified_stat.st_size
content.mtime_ns = verified_stat.st_mtime_ns
else:
split_content(session, content, verified_stat, hash_value=stored_hash)
processed += 1
return processed
def live_contents_under_prefixes(session: Session, prefixes: list[str]) -> list[AssetContent]:
if not prefixes:
return []
return list(
session.scalars(
sa.select(AssetContent).where(
AssetContent.is_missing.is_(False),
sa.or_(*(sql_path_under_prefix(AssetContent.path, prefix) for prefix in prefixes)),
)
)
)