mirror of
https://github.com/Comfy-Org/ComfyUI.git
synced 2026-10-02 02:48:11 -05:00
* 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.
218 lines
7.7 KiB
Python
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)),
|
|
)
|
|
)
|
|
)
|