mirror of
https://github.com/Comfy-Org/ComfyUI.git
synced 2026-10-02 19:08:02 -05:00
* review-stack 1/4: code (37 files, +3217/-3958) Review-and-land stack for synap5e/feat/asset-record-content-split, generated by review-stack.py. Once approved, merges DOWN into the layer below (a fast-forward); only the bottom layer squash-merges into the real base. See ~/adocs/review-stack.md. Rule: path not under tests-unit/ or tests/ Question: Is the logic change right? Source tip:7007d18582Merge-base:783545f689* review-stack 2/4: tests-removed (24 files, +274/-8220) Review-and-land stack for synap5e/feat/asset-record-content-split, generated by review-stack.py. Once approved, merges DOWN into the layer below (a fast-forward); only the bottom layer squash-merges into the real base. See ~/adocs/review-stack.md. Rule: test file deleted, or modified with deleted/(added+deleted) >= 0.9 Question: For each dropped assertion: obsolete by a ruling, or covered by a tests-new test? Source tip:7007d18582Merge-base:783545f689* review-stack 3/4: tests-changed (13 files, +1043/-1218) Review-and-land stack for synap5e/feat/asset-record-content-split, generated by review-stack.py. Once approved, merges DOWN into the layer below (a fast-forward); only the bottom layer squash-merges into the real base. See ~/adocs/review-stack.md. Rule: remaining modified test files (incl. conftest.py / helpers) Question: Did the edits weaken an existing check? Source tip:7007d18582Merge-base:783545f689* review-stack 4/4: tests-new (46 files, +8601/-0) Review-and-land stack for synap5e/feat/asset-record-content-split, generated by review-stack.py. Once approved, merges DOWN into the layer below (a fast-forward); only the bottom layer squash-merges into the real base. See ~/adocs/review-stack.md. Rule: test file added Question: Is the code layer well covered? Source tip:7007d18582Merge-base:783545f689* review-stack 5/6: code (13 files, +351/-104) Review-and-land stack for synap5e/feat/assets-di, generated by review-stack.py. Once approved, merges DOWN into the layer below (a fast-forward); only the bottom layer squash-merges into the real base. See ~/adocs/review-stack.md. Rule: path not under tests-unit/ or tests/ Question: Is the logic change right? Source tip:eca2c74bffMerge-base:20d59d2a5f* review-stack 6/6: tests (8 files, +753/-238) Review-and-land stack for synap5e/feat/assets-di, generated by review-stack.py. Once approved, merges DOWN into the layer below (a fast-forward); only the bottom layer squash-merges into the real base. See ~/adocs/review-stack.md. Rule: every changed file under tests-unit/ or tests/ (added, modified, or deleted) Question: Is the code layer well covered, and did any edit weaken an existing check? Source tip:eca2c74bffMerge-base:20d59d2a5f* review-stack 7/8: ported-fixes (42 files, +1361/-180) Review-and-land stack for synap5e/feat/assets-di-v2, generated by review-stack.py conventions (hand-built continuation layer; see the PR body). Once approved, merges DOWN into the layer below (a fast-forward); only the bottom layer squash-merges into the real base. See ~/adocs/review-stack.md. Rule: the 11 base-branch fix/docs commits 595cd6e4..94d7185b cherry-picked across the DI refactor (7efdd1d7excluded, superseded by layer 8) Question: was each base fix ported faithfully across the DI refactor? Source tip: 6841881069284803b902b4a9e33bdcda13126771 Merge-base:7fdfb40f4b* review-stack 8/8: defensive-parity (4 files, +36/-3) Review-and-land stack for synap5e/feat/assets-di-v2, generated by review-stack.py conventions (hand-built continuation layer; see the PR body). Once approved, merges DOWN into the layer below (a fast-forward); only the bottom layer squash-merges into the real base. See ~/adocs/review-stack.md. Rule: match-or-improve master's dependency defenses — NoAssets selection when DB deps unavailable (7efdd1d7's outcome via the DI seam), requirements warning before assets imports, blake3 in the guarded dependency set Question: does each degradation path now match or improve master's behavior? Source tip:ebc2cfeebcMerge-base:7fdfb40f4b* fix(assets): only discard content rows this operation actually inserted CR-9: Enumerated all six create_content call sites. Only scanner seeding and the three ingest registration paths track IDs for failure cleanup. * fix(assets): reject hash-only uploads with FEATURE_DISABLED when hashing is off CodeRabbit finding CR-2: reject hash-only multipart uploads before create_from_hash when hashing is disabled. * fix(assets): seed persists the stat it verified CR-7: persist the fresh seed-time restat instead of walk-time spec values. * fix(assets): route database lock failures to the lock guidance CR-16: route file-lock startup failures through the existing lock guidance and exit path. * fix(assets): drop the inaccurate temp-cleanup claim from the shutdown warning References CR-10. * fix(assets): walk the output root after execution so undeclared outputs register promptly Custom nodes that write files into the output directory without declaring them in output_ui only became assets when the next full walk happened - a frontend GET /object_info or a restart. Headless and API-only sessions never trigger either, so those files never converged into the asset database. The post-execution hook now requests a FULL scan of the output root instead of an enrich-only pass. The seeder's pending-request queue was generalised from enrich-specific to carrying a scan phase, so the request starts immediately when the seeder is idle and coalesces (escalating to FULL on a phase mismatch) when a scan is already running. queue_output_enrichment is renamed to queue_output_scan across the protocol, the NoAssets no-op and the call site. References FIX-6. * chore(assets): remove seeder paths orphaned by the output-scan change 45c2f96e rerouted both former enrich call sites to start()/enqueue_scan(), leaving two seeder methods that look live but are not. Review round F2 raised this along with four smaller items; the user's disposition was to fix all six here. - Delete start_enrich: zero callers repo-wide after 45c2f96e. - Delete enqueue_enrich: no production callers; its ~18 call sites in tests/test_asset_seeder.py move to enqueue_scan(phase=ScanPhase.ENRICH) with their semantics unchanged. The deletion forces the half-done class renames (TestEnqueueEnrich* -> TestEnqueueScan*, consistent with the already-renamed TestPendingScanDrain) and restores the module docstring that was dropped rather than reworded. - Document at manager.queue_output_scan that ScanPhase.FULL per debounce window is the deliberate, user-ratified trade, so it is not optimised back to ENRICH without revisiting the decision. - Document that SeedAssetSpec.size_bytes/mtime_ns are walk-time diagnostics only - production persists the seed-time restat since CR-7. - Export create_content_reporting_insert from the queries facade and fold scanner.py's direct-module import into the existing facade block. - Harden test_queue_output_scan_does_not_duplicate_declared_output against a vacuous pass: it now asserts the seeder finished without errors and that an undeclared sibling written into the same directory WAS registered by the same scan, proving the walk actually ran. No production behaviour changes beyond the two deletions. References F2-cleanup. * chore: comment cleanup Comment-Gate: 18 quarantined * fix(assets): preserve pause across the seeder's pending-scan drain pause() runs before every prompt, while pending-scan enqueue and resume only run inside the debounced gc-interval gate. If the active scan finishes just after the next prompt's pause, its finally block resets the seeder to idle and the pending drain starts a replacement with the run gate open, so resume becomes a no-op. Capture pausedness under the lock before resetting to idle, then start the drained scan already paused. Setting the state and gate before launching the thread avoids the start-then-reclear window and lets resume release the existing scan checkpoints. Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai> * test(assets): pin job_id absence for scan-discovered assets Owner ruling, recorded 2026-09-03 in the stack-9-hardening planning notepad: scan-discovered assets — including undeclared outputs found by the post-execution walk — carry job_id = None, always; only emission-time registration (output_ui declaration) attributes a job; attributing walk finds to the most recent prompt would be a temporal-correlation guess that is wrong exactly when prompts interleave; None is honest provenance. Do NOT add proximity-based attribution heuristics to the scanner. Ratified against Jacob Segal's cross-job-attribution concern (2026-09-08 review meeting) — a wrongly-attributed asset could mean one user's cloud job sees another user's asset. Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai> * [review-stack 10/10] assets-tests (#16218) * test(execution): run the battery with assets enabled and assert asset-system health at teardown * test(execution): cover list-shaped outputs registering assets Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai> --------- Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai> * [review-stack 11/11] review-fixes (#16261) * fix(assets): only exit on database file-lock timeout when assets are enabled * test(assets): pin live_contents_under_prefixes path-filtering semantics * perf(assets): push live-content prefix filtering into SQL * test(assets): declare per-entry intent in the path-prefix corpus * test(assets): normalize POSIX-literal path expectations for Windows * test(assets): force observable stat changes and close-before-mutate on Windows-sensitive rewrites * test(assets): force an observable mtime change in the hash-mode split test --------- Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai> Co-authored-by: guill <jacob.e.segal@gmail.com>
601 lines
20 KiB
Python
601 lines
20 KiB
Python
"""Walks the asset roots and turns what it finds into catalog rows: collecting
|
|
paths, building specs, seeding new content and records, then enriching them
|
|
with metadata and hashes. Each spec is seeded inside its own savepoint, so one
|
|
file whose row conflicts cannot discard the work done for the files around it.
|
|
Enrichment counts as progress only when it produced what was asked of it — a
|
|
requested hash that could not be computed is no progress, which is what bounds
|
|
a pass over a file the server cannot read.
|
|
"""
|
|
|
|
import logging
|
|
import os
|
|
from dataclasses import dataclass
|
|
from pathlib import Path
|
|
from typing import Callable, Literal, TypedDict
|
|
|
|
import folder_paths
|
|
import sqlalchemy as sa
|
|
from sqlalchemy.exc import IntegrityError
|
|
from sqlalchemy.orm import Session
|
|
from app.assets import mode
|
|
from app.assets.database.queries import (
|
|
create_content_reporting_insert,
|
|
mark_content_missing,
|
|
create_record,
|
|
)
|
|
from app.assets.database.models import Asset, AssetContent
|
|
from app.assets.helpers import sql_path_under_prefix, to_stored_hash
|
|
from app.assets.lifecycle import get_excluded_scan_roots
|
|
from app.assets.scanner_changes import (
|
|
clear_pending_verifications,
|
|
detect_content_change,
|
|
drain_pending_verifications,
|
|
is_path_under_prefixes,
|
|
live_contents_under_prefixes,
|
|
pending_recovery_count,
|
|
recover_missing_content,
|
|
)
|
|
from app.assets.scanner_admission import (
|
|
PARTIAL_DOWNLOAD_EXTENSIONS as PARTIAL_DOWNLOAD_EXTENSIONS,
|
|
_WATCH_LIST as _WATCH_LIST,
|
|
_WatchEntry as _WatchEntry,
|
|
_should_skip_extension,
|
|
_two_stat_admit,
|
|
tick_watch_list as tick_watch_list,
|
|
)
|
|
from app.assets.services.file_utils import get_mtime_ns, is_visible, list_files_recursively
|
|
from app.assets.services.image_dimensions import extract_image_dimensions
|
|
from app.assets.services.metadata_extract import ExtractedMetadata, extract_file_metadata
|
|
from app.assets.services.path_utils import (
|
|
compute_loader_path,
|
|
get_comfy_models_folders,
|
|
get_name_and_tags_from_asset_path,
|
|
)
|
|
from app.assets.services.ingest import _discard_unreferenced_content
|
|
from app.assets.services.snapshot_hash import snapshot_hash
|
|
from app.database.db import create_session
|
|
|
|
__all__ = [
|
|
"clear_pending_verifications",
|
|
"drain_pending_verifications",
|
|
"pending_recovery_count",
|
|
]
|
|
|
|
|
|
# Temp is deliberately absent: it is wiped before every scan, so walking it finds nothing.
|
|
RootType = Literal["models", "input", "output"]
|
|
|
|
|
|
class SeedAssetSpec(TypedDict):
|
|
|
|
abs_path: str
|
|
# Walk-time diagnostics only: seeding persists the seed-time restat instead.
|
|
size_bytes: int
|
|
mtime_ns: int
|
|
info_name: str
|
|
tags: list[str]
|
|
fname: str | None
|
|
metadata: ExtractedMetadata | None
|
|
mime_type: str | None
|
|
job_id: str | None
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class UnenrichedContent:
|
|
content_id: str
|
|
record_id: str
|
|
file_path: str
|
|
|
|
|
|
def _log_scan_error(phase: str, error: OSError) -> None:
|
|
error_type = (
|
|
"permission_denied" if isinstance(error, PermissionError) else "os_error"
|
|
)
|
|
logging.warning("Asset scan error: phase=%s error_type=%s", phase, error_type)
|
|
|
|
|
|
def get_scan_prefixes_for_root(root: RootType) -> list[str]:
|
|
if root == "models":
|
|
bases: list[str] = []
|
|
for _bucket, paths, _exts in get_comfy_models_folders():
|
|
bases.extend(paths)
|
|
return [os.path.abspath(p) for p in bases]
|
|
if root == "input":
|
|
return [os.path.abspath(folder_paths.get_input_directory())]
|
|
if root == "output":
|
|
return [os.path.abspath(folder_paths.get_output_directory())]
|
|
return []
|
|
|
|
|
|
def get_owned_prefixes() -> list[str]:
|
|
"""Every directory an asset may live in; references outside these are marked missing."""
|
|
scan_roots: tuple[RootType, ...] = ("models", "input", "output")
|
|
prefixes = [p for root in scan_roots for p in get_scan_prefixes_for_root(root)]
|
|
return prefixes + get_temp_prefixes()
|
|
|
|
|
|
def get_temp_prefixes() -> list[str]:
|
|
temp_dir = os.path.abspath(folder_paths.get_temp_directory())
|
|
if temp_dir in get_excluded_scan_roots():
|
|
return []
|
|
return [temp_dir]
|
|
|
|
|
|
def collect_models_files() -> list[str]:
|
|
out: list[str] = []
|
|
for folder_name, bases, _exts in get_comfy_models_folders():
|
|
rel_files = folder_paths.get_filename_list(folder_name) or []
|
|
for rel_path in rel_files:
|
|
if not all(is_visible(part) for part in Path(rel_path).parts):
|
|
continue
|
|
abs_path = folder_paths.get_full_path(folder_name, rel_path)
|
|
if not abs_path:
|
|
continue
|
|
abs_path = os.path.abspath(abs_path)
|
|
allowed = False
|
|
abs_p = Path(abs_path)
|
|
for b in bases:
|
|
if abs_p.is_relative_to(os.path.abspath(b)):
|
|
allowed = True
|
|
break
|
|
if allowed:
|
|
out.append(abs_path)
|
|
return out
|
|
|
|
|
|
def sync_references_with_filesystem(
|
|
session,
|
|
root: RootType,
|
|
collect_existing_paths: bool = False,
|
|
) -> set[str] | None:
|
|
return sync_prefixes_with_filesystem(
|
|
session,
|
|
get_scan_prefixes_for_root(root),
|
|
collect_existing_paths=collect_existing_paths,
|
|
)
|
|
|
|
|
|
def sync_prefixes_with_filesystem(
|
|
session: Session,
|
|
prefixes: list[str],
|
|
collect_existing_paths: bool = False,
|
|
) -> set[str] | None:
|
|
if not prefixes:
|
|
return set() if collect_existing_paths else None
|
|
|
|
survivors: set[str] = set()
|
|
for content in live_contents_under_prefixes(session, prefixes):
|
|
try:
|
|
stat_result = os.stat(content.path, follow_symlinks=True)
|
|
except FileNotFoundError:
|
|
mark_content_missing(session, content.id)
|
|
except PermissionError as e:
|
|
_log_scan_error("reference_stat", e)
|
|
logging.debug("Permission denied accessing %s", content.path)
|
|
except OSError as e:
|
|
_log_scan_error("reference_stat", e)
|
|
logging.debug("OSError checking %s: %s", content.path, e)
|
|
mark_content_missing(session, content.id)
|
|
else:
|
|
detect_content_change(
|
|
session,
|
|
content,
|
|
stat_result,
|
|
hashing_is_enabled=mode.hashing_enabled(),
|
|
)
|
|
survivors.add(os.path.abspath(content.path))
|
|
|
|
return survivors if collect_existing_paths else None
|
|
|
|
|
|
def _is_under_prefixes(path: str, prefixes: list[str]) -> bool:
|
|
return is_path_under_prefixes(path, prefixes)
|
|
|
|
|
|
def sync_root_safely(root: RootType) -> set[str]:
|
|
"""Sync a single root's references with the filesystem.
|
|
|
|
Returns survivors (existing paths) or empty set on failure.
|
|
"""
|
|
try:
|
|
with create_session() as sess:
|
|
survivors = sync_references_with_filesystem(
|
|
sess,
|
|
root,
|
|
collect_existing_paths=True,
|
|
)
|
|
sess.commit()
|
|
return survivors or set()
|
|
except Exception as e:
|
|
logging.exception("fast DB scan failed for %s: %s", root, e)
|
|
return set()
|
|
|
|
|
|
def sync_temp_references_safely() -> None:
|
|
"""Retire temp references whose file is gone; temp is never scanned, so nothing else stats them."""
|
|
try:
|
|
with create_session() as sess:
|
|
sync_prefixes_with_filesystem(sess, get_temp_prefixes())
|
|
sess.commit()
|
|
except Exception as e:
|
|
logging.exception("temp reference sync failed: %s", e)
|
|
|
|
|
|
def mark_missing_outside_prefixes_safely(prefixes: list[str]) -> int:
|
|
"""Mark references as missing when outside the given prefixes.
|
|
|
|
This is a non-destructive soft-delete. Returns count marked or 0 on failure.
|
|
"""
|
|
try:
|
|
with create_session() as sess:
|
|
count = mark_contents_missing_outside_prefixes(sess, prefixes)
|
|
sess.commit()
|
|
return count
|
|
except Exception as e:
|
|
logging.exception("marking missing assets failed: %s", e)
|
|
return 0
|
|
|
|
|
|
def mark_contents_missing_outside_prefixes(
|
|
session: Session, prefixes: list[str]
|
|
) -> int:
|
|
contents = session.scalars(
|
|
sa.select(AssetContent).where(AssetContent.is_missing.is_(False))
|
|
)
|
|
missing = [content for content in contents if not _is_under_prefixes(content.path, prefixes)]
|
|
for content in missing:
|
|
mark_content_missing(session, content.id)
|
|
return len(missing)
|
|
|
|
|
|
def collect_paths_for_roots(roots: tuple[RootType, ...]) -> list[str]:
|
|
"""Collect all file paths for the given roots."""
|
|
paths: list[str] = []
|
|
if "models" in roots:
|
|
paths.extend(collect_models_files())
|
|
if "input" in roots:
|
|
paths.extend(list_files_recursively(folder_paths.get_input_directory()))
|
|
if "output" in roots:
|
|
paths.extend(list_files_recursively(folder_paths.get_output_directory()))
|
|
return paths
|
|
|
|
|
|
def build_asset_specs(
|
|
paths: list[str],
|
|
existing_paths: set[str],
|
|
enable_metadata_extraction: bool = True,
|
|
) -> tuple[list[SeedAssetSpec], set[str], int]:
|
|
"""Build asset specs from paths, returning (specs, tag_pool, skipped_count).
|
|
|
|
Args:
|
|
paths: List of file paths to process
|
|
existing_paths: Set of paths that already exist in the database
|
|
enable_metadata_extraction: If True, extract tier 1 & 2 metadata
|
|
"""
|
|
specs: list[SeedAssetSpec] = []
|
|
tag_pool: set[str] = set()
|
|
skipped = 0
|
|
candidates: list[tuple[str, os.stat_result]] = []
|
|
|
|
for p in paths:
|
|
abs_p = os.path.abspath(p)
|
|
if _should_skip_extension(abs_p):
|
|
skipped += 1
|
|
continue
|
|
if abs_p in existing_paths:
|
|
skipped += 1
|
|
continue
|
|
try:
|
|
stat_p = os.stat(abs_p, follow_symlinks=True)
|
|
except FileNotFoundError:
|
|
continue
|
|
except OSError as e:
|
|
_log_scan_error("discovery_stat", e)
|
|
continue
|
|
if not stat_p.st_size:
|
|
continue
|
|
candidates.append((abs_p, stat_p))
|
|
|
|
admitted_paths, _ = _two_stat_admit(candidates)
|
|
candidate_stats = dict(candidates)
|
|
for abs_p in admitted_paths:
|
|
stat_p = candidate_stats[abs_p]
|
|
name, tags = get_name_and_tags_from_asset_path(abs_p)
|
|
rel_fname = compute_loader_path(abs_p)
|
|
|
|
# Extract metadata (tier 1: filesystem, tier 2: safetensors header)
|
|
metadata = None
|
|
if enable_metadata_extraction:
|
|
metadata = extract_file_metadata(
|
|
abs_p,
|
|
stat_result=stat_p,
|
|
relative_filename=rel_fname,
|
|
)
|
|
|
|
mime_type = metadata.content_type if metadata else None
|
|
specs.append(
|
|
{
|
|
"abs_path": abs_p,
|
|
"size_bytes": stat_p.st_size,
|
|
"mtime_ns": get_mtime_ns(stat_p),
|
|
"info_name": name,
|
|
"tags": tags,
|
|
"fname": rel_fname,
|
|
"metadata": metadata,
|
|
"mime_type": mime_type,
|
|
"job_id": None,
|
|
}
|
|
)
|
|
tag_pool.update(tags)
|
|
|
|
return specs, tag_pool, skipped
|
|
|
|
|
|
def seed_asset_specs(session: Session, specs: list[SeedAssetSpec]) -> int:
|
|
created = 0
|
|
created_content_ids: list[str] = []
|
|
try:
|
|
for spec in specs:
|
|
path = os.path.abspath(spec["abs_path"])
|
|
try:
|
|
with session.begin_nested():
|
|
try:
|
|
stat_result = os.stat(path, follow_symlinks=True)
|
|
except OSError:
|
|
logging.warning("Skipping vanished asset during scan: %s", path)
|
|
continue
|
|
try:
|
|
recovery = recover_missing_content(
|
|
session,
|
|
path,
|
|
stat_result,
|
|
hashing_is_enabled=mode.hashing_enabled(),
|
|
)
|
|
except OSError:
|
|
logging.warning("Skipping vanished asset during scan: %s", path)
|
|
continue
|
|
if recovery != "no_match":
|
|
continue
|
|
content, inserted = create_content_reporting_insert(
|
|
session,
|
|
path=path,
|
|
hash=None,
|
|
size_bytes=stat_result.st_size,
|
|
mtime_ns=get_mtime_ns(stat_result),
|
|
)
|
|
if inserted:
|
|
created_content_ids.append(content.id)
|
|
existing_record = session.scalar(
|
|
sa.select(Asset.id).where(Asset.content_id == content.id).limit(1)
|
|
)
|
|
if existing_record is not None:
|
|
continue
|
|
create_record(
|
|
session,
|
|
content_id=content.id,
|
|
name=spec["info_name"],
|
|
mime_type=spec["mime_type"],
|
|
job_id=spec["job_id"],
|
|
loader_path=spec["fname"],
|
|
tags=spec["tags"],
|
|
)
|
|
created += 1
|
|
except IntegrityError:
|
|
logging.warning("Skipping asset whose row conflicts during scan: %s", path)
|
|
continue
|
|
except Exception:
|
|
session.rollback()
|
|
for content_id in created_content_ids:
|
|
_discard_unreferenced_content(session, content_id)
|
|
raise
|
|
return created
|
|
|
|
|
|
def insert_asset_specs(specs: list[SeedAssetSpec], _tag_pool: set[str]) -> int:
|
|
if not specs:
|
|
return 0
|
|
with create_session() as sess:
|
|
created = seed_asset_specs(sess, specs)
|
|
sess.commit()
|
|
return created
|
|
|
|
|
|
def get_unenriched_assets_for_roots(
|
|
roots: tuple[RootType, ...],
|
|
compute_hashes: bool,
|
|
limit: int = 1000,
|
|
) -> list[UnenrichedContent]:
|
|
prefixes: list[str] = []
|
|
for root in roots:
|
|
prefixes.extend(get_scan_prefixes_for_root(root))
|
|
|
|
if not prefixes:
|
|
return []
|
|
|
|
with create_session() as sess:
|
|
query = (
|
|
sa.select(AssetContent.id, Asset.id, AssetContent.path)
|
|
.join(Asset, Asset.content_id == AssetContent.id)
|
|
.where(AssetContent.is_missing.is_(False))
|
|
)
|
|
if compute_hashes:
|
|
query = query.where(
|
|
sa.or_(
|
|
AssetContent.hash.is_(None),
|
|
Asset.system_metadata.is_(None),
|
|
)
|
|
)
|
|
else:
|
|
query = query.where(Asset.system_metadata.is_(None))
|
|
query = query.where(
|
|
sa.or_(
|
|
*(sql_path_under_prefix(AssetContent.path, p) for p in prefixes)
|
|
)
|
|
)
|
|
rows = sess.execute(query.order_by(Asset.id).limit(limit)).all()
|
|
|
|
return [
|
|
UnenrichedContent(content_id, record_id, file_path)
|
|
for content_id, record_id, file_path in rows
|
|
]
|
|
|
|
|
|
def enrich_asset(
|
|
session,
|
|
file_path: str,
|
|
content_id: str,
|
|
record_id: str,
|
|
extract_metadata: bool = True,
|
|
compute_hash: bool = False,
|
|
) -> bool:
|
|
"""Enrich a single asset with metadata and/or hash.
|
|
|
|
Args:
|
|
session: Database session (caller manages lifecycle)
|
|
file_path: Absolute path to the file
|
|
content_id: ID of the content to update
|
|
record_id: ID of the record to update
|
|
extract_metadata: If True, extract safetensors header and mime type
|
|
compute_hash: If True, compute blake3 hash
|
|
|
|
Returns:
|
|
Whether enrichment changed the B-schema record or content
|
|
"""
|
|
try:
|
|
stat_p = os.stat(file_path, follow_symlinks=True)
|
|
except FileNotFoundError:
|
|
return False
|
|
except OSError as e:
|
|
_log_scan_error("enrichment_stat", e)
|
|
return False
|
|
|
|
initial_mtime_ns = get_mtime_ns(stat_p)
|
|
rel_fname = compute_loader_path(file_path)
|
|
mime_type: str | None = None
|
|
metadata = None
|
|
|
|
if extract_metadata:
|
|
metadata = extract_file_metadata(
|
|
file_path,
|
|
stat_result=stat_p,
|
|
relative_filename=rel_fname,
|
|
)
|
|
if metadata:
|
|
mime_type = metadata.content_type
|
|
|
|
content = session.get(AssetContent, content_id)
|
|
|
|
digest: str | None = None
|
|
stored_hash: str | None = None
|
|
verified_stat: os.stat_result | None = None
|
|
hash_requested = compute_hash and content is not None and content.hash is None
|
|
if hash_requested:
|
|
try:
|
|
snapshot = snapshot_hash(file_path)
|
|
if snapshot is None:
|
|
logging.warning(
|
|
"File modified during hashing (snapshot unstable), discarding hash: %s",
|
|
file_path,
|
|
)
|
|
return False
|
|
digest, verified_stat = snapshot
|
|
stored_hash = to_stored_hash(digest)
|
|
except Exception as e:
|
|
if isinstance(e, OSError):
|
|
_log_scan_error("hashing", e)
|
|
else:
|
|
logging.warning("Failed to hash %s: %s", file_path, e)
|
|
|
|
record = session.get(Asset, record_id)
|
|
if content is None or record is None or content.mtime_ns != initial_mtime_ns:
|
|
session.rollback()
|
|
logging.info(
|
|
"Content %s mtime changed during enrichment, discarding stale result",
|
|
content_id,
|
|
)
|
|
return False
|
|
|
|
# Non-NULL system_metadata permanently excludes the row from re-enrichment, so a
|
|
# disagreement here must discard the metadata too, not just the hash.
|
|
if verified_stat is not None and (
|
|
get_mtime_ns(verified_stat) != initial_mtime_ns
|
|
or verified_stat.st_size != stat_p.st_size
|
|
):
|
|
session.rollback()
|
|
logging.info(
|
|
"Content %s changed between its metadata read and its hash read, "
|
|
"discarding stale result",
|
|
content_id,
|
|
)
|
|
return False
|
|
|
|
if extract_metadata and metadata:
|
|
system_metadata = metadata.to_user_metadata()
|
|
if mime_type and mime_type.startswith("image/"):
|
|
dims = extract_image_dimensions(file_path, mime_type=mime_type)
|
|
if dims:
|
|
system_metadata.update(dims)
|
|
record.system_metadata = {**(record.system_metadata or {}), **system_metadata}
|
|
|
|
if stored_hash:
|
|
content.hash = stored_hash
|
|
if mime_type:
|
|
record.mime_type = mime_type
|
|
|
|
session.commit()
|
|
|
|
if hash_requested and stored_hash is None:
|
|
return False
|
|
return stored_hash is not None or metadata is not None or mime_type is not None
|
|
|
|
|
|
def enrich_assets_batch(
|
|
rows: list[UnenrichedContent],
|
|
extract_metadata: bool = True,
|
|
compute_hash: bool = False,
|
|
interrupt_check: Callable[[], bool] | None = None,
|
|
) -> tuple[int, list[str]]:
|
|
"""Enrich a batch of assets.
|
|
|
|
Uses a single DB session for the entire batch, committing after each
|
|
individual asset to avoid long-held transactions while eliminating
|
|
per-asset session creation overhead.
|
|
|
|
Args:
|
|
rows: List of UnenrichedReferenceRow from get_unenriched_assets_for_roots
|
|
extract_metadata: If True, extract metadata for each asset
|
|
compute_hash: If True, compute hash for each asset
|
|
interrupt_check: Optional non-blocking callable that returns True if
|
|
the operation should be interrupted (e.g. paused or cancelled)
|
|
|
|
Returns:
|
|
Tuple of (enriched_count, failed_reference_ids)
|
|
"""
|
|
enriched = 0
|
|
failed_ids: list[str] = []
|
|
|
|
with create_session() as sess:
|
|
for row in rows:
|
|
if interrupt_check is not None and interrupt_check():
|
|
break
|
|
|
|
try:
|
|
updated = enrich_asset(
|
|
sess,
|
|
file_path=row.file_path,
|
|
content_id=row.content_id,
|
|
record_id=row.record_id,
|
|
extract_metadata=extract_metadata,
|
|
compute_hash=compute_hash,
|
|
)
|
|
if updated:
|
|
enriched += 1
|
|
else:
|
|
failed_ids.append(row.record_id)
|
|
except Exception as e:
|
|
logging.warning("Failed to enrich %s: %s", row.file_path, e)
|
|
sess.rollback()
|
|
failed_ids.append(row.record_id)
|
|
|
|
return enriched, failed_ids
|