From 69a54f82dc43e4805dade03ea1bc520b6c6c39dc Mon Sep 17 00:00:00 2001 From: Simon Pinfold Date: Thu, 1 Oct 2026 13:07:05 -0700 Subject: [PATCH 01/15] fix(assets): copy uploads into place when the destination is on another volume os.replace cannot rename across volumes (EXDEV; WinError 17 on Windows), so an upload into an input/output folder on a different drive than the temp directory failed. On EXDEV, copy to a hidden staging file beside the destination and rename that into place, so a partial copy is never visible under the final name; the staging file is removed on any failure. The copy's own stat is recorded, since a copy can change the mtime. --- app/assets/services/ingest.py | 46 +++- .../services/test_upload_cross_device.py | 247 ++++++++++++++++++ 2 files changed, 287 insertions(+), 6 deletions(-) create mode 100644 tests-unit/assets_test/services/test_upload_cross_device.py diff --git a/app/assets/services/ingest.py b/app/assets/services/ingest.py index 52bcde4c0fe..d362eb5622b 100644 --- a/app/assets/services/ingest.py +++ b/app/assets/services/ingest.py @@ -1,16 +1,19 @@ """Turns incoming bytes into catalogued assets: multipart uploads moved into a hash-addressed destination, files registered where they already sit, and records created from a hash the catalog already holds. Every path persists the -stat that hashing verified, so a row's recorded size and mtime describe the -same observation as its hash. A live row already at the destination is +stat that hashing verified (an upload copied across volumes, the copy's stat), +so a row's recorded size and mtime describe the same observation as its hash. A live row already at the destination is reconciled before the write, so an upload never adopts a fresh hash onto records created for bytes it just replaced. """ import contextlib +import errno import logging import mimetypes import os +import shutil +import tempfile from typing import Any, NamedTuple, Sequence from sqlalchemy import func, select @@ -184,10 +187,41 @@ def _guess_upload_mime_type( return guessed or "application/octet-stream" -def _move_temp_to_dest(temp_path: str, dest_abs: str) -> None: +def _copy_across_devices(temp_path: str, dest_abs: str, expected_size: int) -> os.stat_result: + """Copy to a hidden sibling of ``dest_abs``, then rename it into place, so a + partial copy is never visible under the final name. Returns the copy's stat.""" + fd, staging = tempfile.mkstemp( + dir=os.path.dirname(dest_abs), prefix=".", suffix=".upload.tmp" + ) + os.close(fd) + try: + shutil.copy2(temp_path, staging) + copied_stat = os.stat(staging) + if copied_stat.st_size != expected_size: + raise OSError(f"copied {copied_stat.st_size} of {expected_size} bytes") + os.replace(staging, dest_abs) + except BaseException: + with contextlib.suppress(OSError): + os.remove(staging) + raise + return copied_stat + + +def _move_temp_to_dest( + temp_path: str, dest_abs: str, verified_stat: os.stat_result +) -> os.stat_result: + """Move the upload into place and return the stat to record for it. A copy + across volumes (EXDEV, also Windows' ERROR_NOT_SAME_DEVICE) can change the + mtime, so that case records the copy's own stat.""" os.makedirs(os.path.dirname(dest_abs), exist_ok=True) try: - os.replace(temp_path, dest_abs) + try: + os.replace(temp_path, dest_abs) + return verified_stat + except OSError as e: + if e.errno != errno.EXDEV: + raise + return _copy_across_devices(temp_path, dest_abs, verified_stat.st_size) except Exception as e: raise RuntimeError(f"failed to move uploaded file into place: {e}") from e @@ -486,10 +520,10 @@ def upload_from_temp_path( content_type = _guess_upload_mime_type( mime_type, client_filename, name, os.path.basename(dest_abs) ) - _move_temp_to_dest(temp_path, dest_abs) + placed_stat = _move_temp_to_dest(temp_path, dest_abs, verified_stat) finally: _remove_temp_path(temp_path) - size_bytes, mtime_ns = verified_stat.st_size, verified_stat.st_mtime_ns + size_bytes, mtime_ns = placed_stat.st_size, placed_stat.st_mtime_ns system_metadata = _extract_system_metadata_sync(dest_abs, content_type) with create_session() as session: _reconcile_live_content_at_path( diff --git a/tests-unit/assets_test/services/test_upload_cross_device.py b/tests-unit/assets_test/services/test_upload_cross_device.py new file mode 100644 index 00000000000..ba77c97397d --- /dev/null +++ b/tests-unit/assets_test/services/test_upload_cross_device.py @@ -0,0 +1,247 @@ +"""Uploads whose destination sits on a different volume than the temp upload +(another drive letter on Windows, another mount on POSIX), where a rename fails +with EXDEV and the bytes have to be copied instead.""" + +import errno +import os +import shutil +import sys +import uuid +from pathlib import Path + +import pytest +from sqlalchemy import select + +import app.assets.mode as mode_module +import app.assets.services.ingest as ingest_module +import folder_paths +from app.assets.database.models import AssetContent +from app.assets.services.ingest import upload_from_temp_path +from app.assets.services.snapshot_hash import snapshot_hash + +_CONTENT = b"cross-device upload bytes" +_SOURCE_MTIME_NS = 1_600_000_000_123_456_789 + + +@pytest.fixture +def hashing_on(): + class FakeArgs: + enable_asset_hashing = True + + mode_module.init(FakeArgs()) + yield + mode_module.init(None) + + +@pytest.fixture +def dirs(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> tuple[Path, Path]: + temp_root = tmp_path / "temp" + input_root = tmp_path / "input" + temp_root.mkdir() + input_root.mkdir() + monkeypatch.setattr(folder_paths, "get_temp_directory", lambda: str(temp_root)) + monkeypatch.setattr(folder_paths, "get_input_directory", lambda: str(input_root)) + return temp_root, input_root + + +def _write_temp(temp_root: Path) -> Path: + unique_dir = temp_root / "uploads" / uuid.uuid4().hex + unique_dir.mkdir(parents=True) + path = unique_dir / ".upload.part" + path.write_bytes(_CONTENT) + os.utime(path, ns=(_SOURCE_MTIME_NS, _SOURCE_MTIME_NS)) + return path + + +def _dest_for(input_root: Path, temp: Path) -> Path: + snapshot = snapshot_hash(str(temp)) + assert snapshot is not None + return input_root / f"{snapshot[0]}.png" + + +def _upload(temp: Path): + return upload_from_temp_path( + temp_path=str(temp), name="photo.png", tags=["input"], client_filename="photo.png" + ) + + +def _fail_first_replace(monkeypatch: pytest.MonkeyPatch, exc: OSError, then=None) -> list: + """Make the first os.replace (temp -> destination) raise ``exc``; later calls + run ``then`` if given, else the real os.replace.""" + real_replace = os.replace + calls: list = [] + + def fake_replace(src, dst): + calls.append((src, dst)) + if len(calls) == 1: + raise exc + if then is not None: + return then(src, dst) + return real_replace(src, dst) + + monkeypatch.setattr(ingest_module.os, "replace", fake_replace) + return calls + + +def _exdev() -> OSError: + return OSError(errno.EXDEV, "Invalid cross-device link") + + +def _assert_nothing_left(temp: Path, input_root: Path, dest: Path | None) -> None: + assert not temp.exists() + assert not temp.parent.exists() + leftovers = sorted(p.name for p in input_root.iterdir()) + assert leftovers == ([dest.name] if dest is not None else []) + + +def test_cross_device_upload_copies_into_place_and_records_the_copys_stat( + mock_create_session, hashing_on, dirs, monkeypatch +): + temp_root, input_root = dirs + temp = _write_temp(temp_root) + dest = _dest_for(input_root, temp) + copied_mtime_ns = _SOURCE_MTIME_NS + 7_000_000_000 + real_copy2 = shutil.copy2 + + def copy2_on_coarser_clock(src, dst): + real_copy2(src, dst) + os.utime(dst, ns=(copied_mtime_ns, copied_mtime_ns)) + + monkeypatch.setattr(ingest_module.shutil, "copy2", copy2_on_coarser_clock) + calls = _fail_first_replace(monkeypatch, _exdev()) + + result = _upload(temp) + + assert result.created_new is True + assert dest.read_bytes() == _CONTENT + staging_src, staging_dst = calls[1] + assert Path(staging_dst) == dest + assert Path(staging_src).parent == input_root + assert Path(staging_src).name.startswith(".") + assert Path(staging_src).name.endswith(".tmp") + _assert_nothing_left(temp, input_root, dest) + with mock_create_session() as session: + content = session.scalars(select(AssetContent)).one() + assert content.path == str(dest) + assert content.size_bytes == len(_CONTENT) + assert content.mtime_ns == dest.stat().st_mtime_ns == copied_mtime_ns + + +def test_copy_failure_leaves_no_destination_and_no_staging_file( + mock_create_session, hashing_on, dirs, monkeypatch +): + temp_root, input_root = dirs + temp = _write_temp(temp_root) + + def partial_copy(src, dst): + Path(dst).write_bytes(_CONTENT[:5]) + raise OSError(errno.ENOSPC, "No space left on device") + + monkeypatch.setattr(ingest_module.shutil, "copy2", partial_copy) + _fail_first_replace(monkeypatch, _exdev()) + + with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): + _upload(temp) + + _assert_nothing_left(temp, input_root, None) + + +def test_short_copy_is_rejected_before_it_takes_the_final_name( + mock_create_session, hashing_on, dirs, monkeypatch +): + temp_root, input_root = dirs + temp = _write_temp(temp_root) + monkeypatch.setattr( + ingest_module.shutil, "copy2", lambda src, dst: Path(dst).write_bytes(_CONTENT[:5]) + ) + _fail_first_replace(monkeypatch, _exdev()) + + with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): + _upload(temp) + + _assert_nothing_left(temp, input_root, None) + + +def test_final_rename_failure_removes_the_staging_file( + mock_create_session, hashing_on, dirs, monkeypatch +): + temp_root, input_root = dirs + temp = _write_temp(temp_root) + + def locked(src, dst): + raise PermissionError(errno.EACCES, "The process cannot access the file") + + _fail_first_replace(monkeypatch, _exdev(), then=locked) + + with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): + _upload(temp) + + _assert_nothing_left(temp, input_root, None) + + +def test_other_move_errors_do_not_fall_back_to_a_copy( + mock_create_session, hashing_on, dirs, monkeypatch +): + temp_root, input_root = dirs + temp = _write_temp(temp_root) + copies: list = [] + monkeypatch.setattr(ingest_module.shutil, "copy2", lambda *a: copies.append(a)) + _fail_first_replace(monkeypatch, PermissionError(errno.EACCES, "Access is denied")) + + with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): + _upload(temp) + + assert copies == [] + _assert_nothing_left(temp, input_root, None) + + +def test_same_device_upload_is_a_plain_rename( + mock_create_session, hashing_on, dirs, monkeypatch +): + temp_root, input_root = dirs + temp = _write_temp(temp_root) + dest = _dest_for(input_root, temp) + copies: list = [] + monkeypatch.setattr(ingest_module.shutil, "copy2", lambda *a: copies.append(a)) + + _upload(temp) + + assert copies == [] + assert dest.read_bytes() == _CONTENT + _assert_nothing_left(temp, input_root, dest) + with mock_create_session() as session: + assert session.scalars(select(AssetContent)).one().mtime_ns == _SOURCE_MTIME_NS + + +@pytest.mark.skipif(sys.platform != "win32", reason="Windows error mapping") +def test_windows_not_same_device_error_maps_to_exdev(): + # ERROR_NOT_SAME_DEVICE (17) is what os.replace raises across drive letters. + assert OSError(None, "not same device", None, 17).errno == errno.EXDEV + + +def test_real_cross_filesystem_upload(mock_create_session, hashing_on, tmp_path, monkeypatch): + shm = Path("/dev/shm") + if not shm.is_dir() or not os.access(shm, os.W_OK): + pytest.skip("no writable /dev/shm") + if shm.stat().st_dev == tmp_path.stat().st_dev: + pytest.skip("/dev/shm and tmp_path share a filesystem") + temp_root = shm / f"comfy-cross-device-{uuid.uuid4().hex}" + input_root = tmp_path / "input" + temp_root.mkdir() + input_root.mkdir() + monkeypatch.setattr(folder_paths, "get_temp_directory", lambda: str(temp_root)) + monkeypatch.setattr(folder_paths, "get_input_directory", lambda: str(input_root)) + try: + temp = _write_temp(temp_root) + dest = _dest_for(input_root, temp) + + _upload(temp) + + assert dest.read_bytes() == _CONTENT + _assert_nothing_left(temp, input_root, dest) + with mock_create_session() as session: + assert session.scalars(select(AssetContent)).one().mtime_ns == ( + dest.stat().st_mtime_ns + ) + finally: + shutil.rmtree(temp_root, ignore_errors=True) From 875f8909fc30cc93c9d620639ab9c09eb634c798 Mon Sep 17 00:00:00 2001 From: Simon Pinfold Date: Thu, 1 Oct 2026 13:16:08 -0700 Subject: [PATCH 02/15] fix(assets): don't fail cross-volume uploads on metadata, verify the source shutil.copy2 also copies permission bits, which raises on filesystems that cannot store them (FAT/exFAT on Linux), failing the upload after the bytes were copied. Copy the bytes, then copy the mode best-effort. The copy's stat is what gets recorded, so check that the source still matches the stat hashing verified once the copy is done; otherwise bytes rewritten after hashing would be published under the old hash. --- app/assets/services/ingest.py | 30 +++++--- .../services/test_upload_cross_device.py | 77 ++++++++++++++----- 2 files changed, 80 insertions(+), 27 deletions(-) diff --git a/app/assets/services/ingest.py b/app/assets/services/ingest.py index d362eb5622b..7b87cb3c041 100644 --- a/app/assets/services/ingest.py +++ b/app/assets/services/ingest.py @@ -1,8 +1,9 @@ """Turns incoming bytes into catalogued assets: multipart uploads moved into a hash-addressed destination, files registered where they already sit, and records created from a hash the catalog already holds. Every path persists the -stat that hashing verified (an upload copied across volumes, the copy's stat), -so a row's recorded size and mtime describe the same observation as its hash. A live row already at the destination is +stat that hashing verified (for an upload copied across volumes, the stat of a +copy taken from that verified file), so a row's recorded size and mtime describe +the same observation as its hash. A live row already at the destination is reconciled before the write, so an upload never adopts a fresh hash onto records created for bytes it just replaced. """ @@ -187,18 +188,29 @@ def _guess_upload_mime_type( return guessed or "application/octet-stream" -def _copy_across_devices(temp_path: str, dest_abs: str, expected_size: int) -> os.stat_result: +def _copy_across_devices( + temp_path: str, dest_abs: str, verified_stat: os.stat_result +) -> os.stat_result: """Copy to a hidden sibling of ``dest_abs``, then rename it into place, so a - partial copy is never visible under the final name. Returns the copy's stat.""" + partial copy is never visible under the final name. The source must still + match the stat hashing verified once the copy is done. Returns the copy's + stat.""" fd, staging = tempfile.mkstemp( dir=os.path.dirname(dest_abs), prefix=".", suffix=".upload.tmp" ) - os.close(fd) try: - shutil.copy2(temp_path, staging) + with os.fdopen(fd, "wb") as dst, open(temp_path, "rb") as src: + shutil.copyfileobj(src, dst) + source_stat = os.fstat(src.fileno()) + if (source_stat.st_size, source_stat.st_mtime_ns) != ( + verified_stat.st_size, + verified_stat.st_mtime_ns, + ): + raise OSError("upload file changed after hashing") + # Best effort: mode bits cannot be set on some filesystems (e.g. FAT). + with contextlib.suppress(OSError): + shutil.copymode(temp_path, staging) copied_stat = os.stat(staging) - if copied_stat.st_size != expected_size: - raise OSError(f"copied {copied_stat.st_size} of {expected_size} bytes") os.replace(staging, dest_abs) except BaseException: with contextlib.suppress(OSError): @@ -221,7 +233,7 @@ def _move_temp_to_dest( except OSError as e: if e.errno != errno.EXDEV: raise - return _copy_across_devices(temp_path, dest_abs, verified_stat.st_size) + return _copy_across_devices(temp_path, dest_abs, verified_stat) except Exception as e: raise RuntimeError(f"failed to move uploaded file into place: {e}") from e diff --git a/tests-unit/assets_test/services/test_upload_cross_device.py b/tests-unit/assets_test/services/test_upload_cross_device.py index ba77c97397d..35e8548ad3f 100644 --- a/tests-unit/assets_test/services/test_upload_cross_device.py +++ b/tests-unit/assets_test/services/test_upload_cross_device.py @@ -99,15 +99,8 @@ def test_cross_device_upload_copies_into_place_and_records_the_copys_stat( ): temp_root, input_root = dirs temp = _write_temp(temp_root) + temp.chmod(0o644) dest = _dest_for(input_root, temp) - copied_mtime_ns = _SOURCE_MTIME_NS + 7_000_000_000 - real_copy2 = shutil.copy2 - - def copy2_on_coarser_clock(src, dst): - real_copy2(src, dst) - os.utime(dst, ns=(copied_mtime_ns, copied_mtime_ns)) - - monkeypatch.setattr(ingest_module.shutil, "copy2", copy2_on_coarser_clock) calls = _fail_first_replace(monkeypatch, _exdev()) result = _upload(temp) @@ -119,12 +112,33 @@ def copy2_on_coarser_clock(src, dst): assert Path(staging_src).parent == input_root assert Path(staging_src).name.startswith(".") assert Path(staging_src).name.endswith(".tmp") + if sys.platform != "win32": + assert dest.stat().st_mode & 0o777 == 0o644 _assert_nothing_left(temp, input_root, dest) with mock_create_session() as session: content = session.scalars(select(AssetContent)).one() assert content.path == str(dest) assert content.size_bytes == len(_CONTENT) - assert content.mtime_ns == dest.stat().st_mtime_ns == copied_mtime_ns + assert content.mtime_ns == dest.stat().st_mtime_ns != _SOURCE_MTIME_NS + + +def test_destination_that_rejects_mode_bits_still_accepts_the_upload( + mock_create_session, hashing_on, dirs, monkeypatch +): + temp_root, input_root = dirs + temp = _write_temp(temp_root) + dest = _dest_for(input_root, temp) + + def no_chmod(src, dst): + raise PermissionError(errno.EPERM, "Operation not permitted") + + monkeypatch.setattr(ingest_module.shutil, "copymode", no_chmod) + _fail_first_replace(monkeypatch, _exdev()) + + _upload(temp) + + assert dest.read_bytes() == _CONTENT + _assert_nothing_left(temp, input_root, dest) def test_copy_failure_leaves_no_destination_and_no_staging_file( @@ -134,10 +148,10 @@ def test_copy_failure_leaves_no_destination_and_no_staging_file( temp = _write_temp(temp_root) def partial_copy(src, dst): - Path(dst).write_bytes(_CONTENT[:5]) + dst.write(src.read(5)) raise OSError(errno.ENOSPC, "No space left on device") - monkeypatch.setattr(ingest_module.shutil, "copy2", partial_copy) + monkeypatch.setattr(ingest_module.shutil, "copyfileobj", partial_copy) _fail_first_replace(monkeypatch, _exdev()) with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): @@ -146,17 +160,44 @@ def partial_copy(src, dst): _assert_nothing_left(temp, input_root, None) -def test_short_copy_is_rejected_before_it_takes_the_final_name( +def _rewrite_same_size(temp: Path) -> None: + temp.write_bytes(_CONTENT[::-1]) + os.utime(temp, ns=(_SOURCE_MTIME_NS + 1, _SOURCE_MTIME_NS + 1)) + + +def test_source_changed_after_hashing_is_not_copied( mock_create_session, hashing_on, dirs, monkeypatch ): temp_root, input_root = dirs temp = _write_temp(temp_root) - monkeypatch.setattr( - ingest_module.shutil, "copy2", lambda src, dst: Path(dst).write_bytes(_CONTENT[:5]) - ) + + def rewrite_then_exdev(src, dst): + _rewrite_same_size(temp) + raise _exdev() + + monkeypatch.setattr(ingest_module.os, "replace", rewrite_then_exdev) + + with pytest.raises(RuntimeError, match="changed after hashing"): + _upload(temp) + + _assert_nothing_left(temp, input_root, None) + + +def test_source_changed_during_the_copy_is_not_published( + mock_create_session, hashing_on, dirs, monkeypatch +): + temp_root, input_root = dirs + temp = _write_temp(temp_root) + real_copyfileobj = shutil.copyfileobj + + def copy_then_rewrite(src, dst): + real_copyfileobj(src, dst) + _rewrite_same_size(temp) + + monkeypatch.setattr(ingest_module.shutil, "copyfileobj", copy_then_rewrite) _fail_first_replace(monkeypatch, _exdev()) - with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): + with pytest.raises(RuntimeError, match="changed after hashing"): _upload(temp) _assert_nothing_left(temp, input_root, None) @@ -185,7 +226,7 @@ def test_other_move_errors_do_not_fall_back_to_a_copy( temp_root, input_root = dirs temp = _write_temp(temp_root) copies: list = [] - monkeypatch.setattr(ingest_module.shutil, "copy2", lambda *a: copies.append(a)) + monkeypatch.setattr(ingest_module.shutil, "copyfileobj", lambda *a: copies.append(a)) _fail_first_replace(monkeypatch, PermissionError(errno.EACCES, "Access is denied")) with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): @@ -202,7 +243,7 @@ def test_same_device_upload_is_a_plain_rename( temp = _write_temp(temp_root) dest = _dest_for(input_root, temp) copies: list = [] - monkeypatch.setattr(ingest_module.shutil, "copy2", lambda *a: copies.append(a)) + monkeypatch.setattr(ingest_module.shutil, "copyfileobj", lambda *a: copies.append(a)) _upload(temp) From 8036cd04ab58570720ad95dd839cdfad4cc3bb57 Mon Sep 17 00:00:00 2001 From: Simon Pinfold Date: Thu, 1 Oct 2026 13:49:37 -0700 Subject: [PATCH 03/15] test(assets): use NTFS-representable mtimes in cross-volume upload tests NTFS stores mtimes at 100 ns, so a sub-microsecond source mtime was truncated and a +1 ns rewrite left the mtime unchanged on Windows. --- tests-unit/assets_test/services/test_upload_cross_device.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests-unit/assets_test/services/test_upload_cross_device.py b/tests-unit/assets_test/services/test_upload_cross_device.py index 35e8548ad3f..a97d5437def 100644 --- a/tests-unit/assets_test/services/test_upload_cross_device.py +++ b/tests-unit/assets_test/services/test_upload_cross_device.py @@ -20,7 +20,7 @@ from app.assets.services.snapshot_hash import snapshot_hash _CONTENT = b"cross-device upload bytes" -_SOURCE_MTIME_NS = 1_600_000_000_123_456_789 +_SOURCE_MTIME_NS = 1_600_000_000_000_000_000 # whole seconds: NTFS keeps 100 ns @pytest.fixture @@ -162,7 +162,7 @@ def partial_copy(src, dst): def _rewrite_same_size(temp: Path) -> None: temp.write_bytes(_CONTENT[::-1]) - os.utime(temp, ns=(_SOURCE_MTIME_NS + 1, _SOURCE_MTIME_NS + 1)) + os.utime(temp, ns=(_SOURCE_MTIME_NS + 2_000_000_000,) * 2) def test_source_changed_after_hashing_is_not_copied( From 7b47e07df451b0356f65f120f6a45b1db44c93b0 Mon Sep 17 00:00:00 2001 From: Simon Pinfold Date: Thu, 1 Oct 2026 15:19:17 -0700 Subject: [PATCH 04/15] fix(assets): fail fast when the destination denies writes tempfile.mkstemp on Windows retries a PermissionError for as long as os.access reports the directory writable, which an ACL deny does not change, so an upload into a write-denied folder hung the server. Create the staging file directly with O_EXCL, retrying only on a name collision. --- app/assets/services/ingest.py | 22 +++++-- .../services/test_upload_cross_device.py | 63 ++++++++++++++++++- 2 files changed, 79 insertions(+), 6 deletions(-) diff --git a/app/assets/services/ingest.py b/app/assets/services/ingest.py index 7b87cb3c041..ea5e8cbc477 100644 --- a/app/assets/services/ingest.py +++ b/app/assets/services/ingest.py @@ -13,8 +13,8 @@ import logging import mimetypes import os +import secrets import shutil -import tempfile from typing import Any, NamedTuple, Sequence from sqlalchemy import func, select @@ -188,6 +188,22 @@ def _guess_upload_mime_type( return guessed or "application/octet-stream" +_STAGING_NAME_ATTEMPTS = 16 + + +def _create_staging_file(directory: str) -> tuple[int, str]: + """Not tempfile.mkstemp: on Windows it retries PermissionError until TMP_MAX + when an ACL denies writes, hanging the upload.""" + flags = os.O_CREAT | os.O_EXCL | os.O_WRONLY | getattr(os, "O_BINARY", 0) + for _ in range(_STAGING_NAME_ATTEMPTS): + path = os.path.join(directory, f".{secrets.token_hex(8)}.upload.tmp") + try: + return os.open(path, flags, 0o666), path + except FileExistsError: + continue + raise FileExistsError(f"no free staging name in {directory}") + + def _copy_across_devices( temp_path: str, dest_abs: str, verified_stat: os.stat_result ) -> os.stat_result: @@ -195,9 +211,7 @@ def _copy_across_devices( partial copy is never visible under the final name. The source must still match the stat hashing verified once the copy is done. Returns the copy's stat.""" - fd, staging = tempfile.mkstemp( - dir=os.path.dirname(dest_abs), prefix=".", suffix=".upload.tmp" - ) + fd, staging = _create_staging_file(os.path.dirname(dest_abs)) try: with os.fdopen(fd, "wb") as dst, open(temp_path, "rb") as src: shutil.copyfileobj(src, dst) diff --git a/tests-unit/assets_test/services/test_upload_cross_device.py b/tests-unit/assets_test/services/test_upload_cross_device.py index a97d5437def..ed415e97861 100644 --- a/tests-unit/assets_test/services/test_upload_cross_device.py +++ b/tests-unit/assets_test/services/test_upload_cross_device.py @@ -99,7 +99,7 @@ def test_cross_device_upload_copies_into_place_and_records_the_copys_stat( ): temp_root, input_root = dirs temp = _write_temp(temp_root) - temp.chmod(0o644) + temp.chmod(0o640) dest = _dest_for(input_root, temp) calls = _fail_first_replace(monkeypatch, _exdev()) @@ -113,7 +113,7 @@ def test_cross_device_upload_copies_into_place_and_records_the_copys_stat( assert Path(staging_src).name.startswith(".") assert Path(staging_src).name.endswith(".tmp") if sys.platform != "win32": - assert dest.stat().st_mode & 0o777 == 0o644 + assert dest.stat().st_mode & 0o777 == 0o640 _assert_nothing_left(temp, input_root, dest) with mock_create_session() as session: content = session.scalars(select(AssetContent)).one() @@ -203,6 +203,65 @@ def copy_then_rewrite(src, dst): _assert_nothing_left(temp, input_root, None) +def _deny_staging_writes(monkeypatch: pytest.MonkeyPatch, input_root: Path) -> list: + """An ACL that denies writes in the destination: creating a file there raises + PermissionError while os.access still reports the directory writable. The + first os.replace raises EXDEV and switches os.name to "nt", where + tempfile.mkstemp retries that PermissionError up to TMP_MAX times.""" + real_open = os.open + attempts: list = [] + + def denying_open(path, flags, mode=0o777, *, dir_fd=None): + if os.path.dirname(path) == str(input_root): + attempts.append(path) + if len(attempts) > 50: + raise RuntimeError("staging creation kept retrying") + raise PermissionError(errno.EACCES, "Access is denied", str(path)) + return real_open(path, flags, mode, dir_fd=dir_fd) + + def exdev_as_windows(src, dst): + monkeypatch.setattr(os, "name", "nt") + raise _exdev() + + monkeypatch.setattr(ingest_module.os, "open", denying_open) + monkeypatch.setattr(ingest_module.os, "replace", exdev_as_windows) + return attempts + + +def test_write_denied_destination_fails_fast( + mock_create_session, hashing_on, dirs, monkeypatch +): + temp_root, input_root = dirs + temp = _write_temp(temp_root) + attempts = _deny_staging_writes(monkeypatch, input_root) + + with pytest.raises(RuntimeError, match="failed to move uploaded file into place: .*denied"): + _upload(temp) + + monkeypatch.undo() + assert len(attempts) == 1 + _assert_nothing_left(temp, input_root, None) + + +def test_staging_name_collision_takes_another_name( + mock_create_session, hashing_on, dirs, monkeypatch +): + temp_root, input_root = dirs + temp = _write_temp(temp_root) + dest = _dest_for(input_root, temp) + taken = iter(["0" * 16, "0" * 16, "1" * 16]) + (input_root / f".{'0' * 16}.upload.tmp").write_bytes(b"someone else's") + monkeypatch.setattr(ingest_module.secrets, "token_hex", lambda _n: next(taken)) + _fail_first_replace(monkeypatch, _exdev()) + + _upload(temp) + + assert dest.read_bytes() == _CONTENT + assert sorted(p.name for p in input_root.iterdir()) == sorted( + [dest.name, f".{'0' * 16}.upload.tmp"] + ) + + def test_final_rename_failure_removes_the_staging_file( mock_create_session, hashing_on, dirs, monkeypatch ): From b0867f9591ca51d086f1babde77407060dcfb426 Mon Sep 17 00:00:00 2001 From: Simon Pinfold Date: Thu, 1 Oct 2026 22:10:28 -0700 Subject: [PATCH 05/15] fix(assets): report a source changed during a cross-volume copy as unstable Raise UploadUnstableError, as the hash-time check does, so the API returns UPLOAD_UNSTABLE instead of a generic internal error. --- app/assets/services/ingest.py | 4 +++- tests-unit/assets_test/services/test_upload_cross_device.py | 6 +++--- 2 files changed, 6 insertions(+), 4 deletions(-) diff --git a/app/assets/services/ingest.py b/app/assets/services/ingest.py index ea5e8cbc477..39ac271449a 100644 --- a/app/assets/services/ingest.py +++ b/app/assets/services/ingest.py @@ -220,7 +220,7 @@ def _copy_across_devices( verified_stat.st_size, verified_stat.st_mtime_ns, ): - raise OSError("upload file changed after hashing") + raise UploadUnstableError("upload file changed after hashing") # Best effort: mode bits cannot be set on some filesystems (e.g. FAT). with contextlib.suppress(OSError): shutil.copymode(temp_path, staging) @@ -248,6 +248,8 @@ def _move_temp_to_dest( if e.errno != errno.EXDEV: raise return _copy_across_devices(temp_path, dest_abs, verified_stat) + except UploadUnstableError: + raise except Exception as e: raise RuntimeError(f"failed to move uploaded file into place: {e}") from e diff --git a/tests-unit/assets_test/services/test_upload_cross_device.py b/tests-unit/assets_test/services/test_upload_cross_device.py index ed415e97861..3d6277dbe26 100644 --- a/tests-unit/assets_test/services/test_upload_cross_device.py +++ b/tests-unit/assets_test/services/test_upload_cross_device.py @@ -16,7 +16,7 @@ import app.assets.services.ingest as ingest_module import folder_paths from app.assets.database.models import AssetContent -from app.assets.services.ingest import upload_from_temp_path +from app.assets.services.ingest import UploadUnstableError, upload_from_temp_path from app.assets.services.snapshot_hash import snapshot_hash _CONTENT = b"cross-device upload bytes" @@ -177,7 +177,7 @@ def rewrite_then_exdev(src, dst): monkeypatch.setattr(ingest_module.os, "replace", rewrite_then_exdev) - with pytest.raises(RuntimeError, match="changed after hashing"): + with pytest.raises(UploadUnstableError, match="changed after hashing"): _upload(temp) _assert_nothing_left(temp, input_root, None) @@ -197,7 +197,7 @@ def copy_then_rewrite(src, dst): monkeypatch.setattr(ingest_module.shutil, "copyfileobj", copy_then_rewrite) _fail_first_replace(monkeypatch, _exdev()) - with pytest.raises(RuntimeError, match="changed after hashing"): + with pytest.raises(UploadUnstableError, match="changed after hashing"): _upload(temp) _assert_nothing_left(temp, input_root, None) From bc46bca79cb2ca062a33d0b70b44731ace99e8ce Mon Sep 17 00:00:00 2001 From: Simon Pinfold Date: Fri, 2 Oct 2026 20:15:40 -0700 Subject: [PATCH 06/15] refactor(assets): reduce the cross-volume upload move to shutil.move Same-volume uploads stay a single os.replace. On EXDEV, fall back to shutil.move and record the destination's stat, removing a truncated copy if the move fails. Drop the staging file, mode copy and source re-check along with their tests. --- app/assets/services/ingest.py | 75 ++----- .../services/test_upload_cross_device.py | 204 ++++-------------- 2 files changed, 56 insertions(+), 223 deletions(-) diff --git a/app/assets/services/ingest.py b/app/assets/services/ingest.py index 39ac271449a..20bbb56eac4 100644 --- a/app/assets/services/ingest.py +++ b/app/assets/services/ingest.py @@ -1,11 +1,11 @@ """Turns incoming bytes into catalogued assets: multipart uploads moved into a hash-addressed destination, files registered where they already sit, and records created from a hash the catalog already holds. Every path persists the -stat that hashing verified (for an upload copied across volumes, the stat of a -copy taken from that verified file), so a row's recorded size and mtime describe -the same observation as its hash. A live row already at the destination is -reconciled before the write, so an upload never adopts a fresh hash onto -records created for bytes it just replaced. +stat that hashing verified (for an upload copied across volumes, the stat of the +copy), so a row's recorded size and mtime describe the same observation as its +hash. A live row already at the destination is reconciled before the write, so +an upload never adopts a fresh hash onto records created for bytes it just +replaced. """ import contextlib @@ -13,7 +13,6 @@ import logging import mimetypes import os -import secrets import shutil from typing import Any, NamedTuple, Sequence @@ -188,57 +187,12 @@ def _guess_upload_mime_type( return guessed or "application/octet-stream" -_STAGING_NAME_ATTEMPTS = 16 - - -def _create_staging_file(directory: str) -> tuple[int, str]: - """Not tempfile.mkstemp: on Windows it retries PermissionError until TMP_MAX - when an ACL denies writes, hanging the upload.""" - flags = os.O_CREAT | os.O_EXCL | os.O_WRONLY | getattr(os, "O_BINARY", 0) - for _ in range(_STAGING_NAME_ATTEMPTS): - path = os.path.join(directory, f".{secrets.token_hex(8)}.upload.tmp") - try: - return os.open(path, flags, 0o666), path - except FileExistsError: - continue - raise FileExistsError(f"no free staging name in {directory}") - - -def _copy_across_devices( - temp_path: str, dest_abs: str, verified_stat: os.stat_result -) -> os.stat_result: - """Copy to a hidden sibling of ``dest_abs``, then rename it into place, so a - partial copy is never visible under the final name. The source must still - match the stat hashing verified once the copy is done. Returns the copy's - stat.""" - fd, staging = _create_staging_file(os.path.dirname(dest_abs)) - try: - with os.fdopen(fd, "wb") as dst, open(temp_path, "rb") as src: - shutil.copyfileobj(src, dst) - source_stat = os.fstat(src.fileno()) - if (source_stat.st_size, source_stat.st_mtime_ns) != ( - verified_stat.st_size, - verified_stat.st_mtime_ns, - ): - raise UploadUnstableError("upload file changed after hashing") - # Best effort: mode bits cannot be set on some filesystems (e.g. FAT). - with contextlib.suppress(OSError): - shutil.copymode(temp_path, staging) - copied_stat = os.stat(staging) - os.replace(staging, dest_abs) - except BaseException: - with contextlib.suppress(OSError): - os.remove(staging) - raise - return copied_stat - - def _move_temp_to_dest( temp_path: str, dest_abs: str, verified_stat: os.stat_result ) -> os.stat_result: - """Move the upload into place and return the stat to record for it. A copy - across volumes (EXDEV, also Windows' ERROR_NOT_SAME_DEVICE) can change the - mtime, so that case records the copy's own stat.""" + """Move the upload into place and return the stat to record for it. Across + volumes (EXDEV, also Windows' ERROR_NOT_SAME_DEVICE) the move is a copy with + its own mtime, so that case records the destination's stat.""" os.makedirs(os.path.dirname(dest_abs), exist_ok=True) try: try: @@ -247,9 +201,16 @@ def _move_temp_to_dest( except OSError as e: if e.errno != errno.EXDEV: raise - return _copy_across_devices(temp_path, dest_abs, verified_stat) - except UploadUnstableError: - raise + dest_existed = os.path.exists(dest_abs) + try: + shutil.move(temp_path, dest_abs) + except BaseException: + # A failed copy (e.g. a full disk) leaves a truncated file at dest. + if not dest_existed: + with contextlib.suppress(OSError): + os.remove(dest_abs) + raise + return os.stat(dest_abs) except Exception as e: raise RuntimeError(f"failed to move uploaded file into place: {e}") from e diff --git a/tests-unit/assets_test/services/test_upload_cross_device.py b/tests-unit/assets_test/services/test_upload_cross_device.py index 3d6277dbe26..52deb62839a 100644 --- a/tests-unit/assets_test/services/test_upload_cross_device.py +++ b/tests-unit/assets_test/services/test_upload_cross_device.py @@ -16,7 +16,7 @@ import app.assets.services.ingest as ingest_module import folder_paths from app.assets.database.models import AssetContent -from app.assets.services.ingest import UploadUnstableError, upload_from_temp_path +from app.assets.services.ingest import upload_from_temp_path from app.assets.services.snapshot_hash import snapshot_hash _CONTENT = b"cross-device upload bytes" @@ -65,26 +65,16 @@ def _upload(temp: Path): ) -def _fail_first_replace(monkeypatch: pytest.MonkeyPatch, exc: OSError, then=None) -> list: - """Make the first os.replace (temp -> destination) raise ``exc``; later calls - run ``then`` if given, else the real os.replace.""" - real_replace = os.replace - calls: list = [] +def _cross_device(monkeypatch: pytest.MonkeyPatch, exc: OSError | None = None) -> None: + """Make renames fail as they do across volumes: os.replace raises ``exc`` + (EXDEV by default), and so does the os.rename shutil.move tries first.""" + exc = exc or OSError(errno.EXDEV, "Invalid cross-device link") - def fake_replace(src, dst): - calls.append((src, dst)) - if len(calls) == 1: - raise exc - if then is not None: - return then(src, dst) - return real_replace(src, dst) + def fail(src, dst): + raise exc - monkeypatch.setattr(ingest_module.os, "replace", fake_replace) - return calls - - -def _exdev() -> OSError: - return OSError(errno.EXDEV, "Invalid cross-device link") + monkeypatch.setattr(ingest_module.os, "replace", fail) + monkeypatch.setattr(ingest_module.os, "rename", fail) def _assert_nothing_left(temp: Path, input_root: Path, dest: Path | None) -> None: @@ -94,65 +84,46 @@ def _assert_nothing_left(temp: Path, input_root: Path, dest: Path | None) -> Non assert leftovers == ([dest.name] if dest is not None else []) -def test_cross_device_upload_copies_into_place_and_records_the_copys_stat( +def test_cross_device_upload_is_copied_into_place( mock_create_session, hashing_on, dirs, monkeypatch ): temp_root, input_root = dirs temp = _write_temp(temp_root) - temp.chmod(0o640) dest = _dest_for(input_root, temp) - calls = _fail_first_replace(monkeypatch, _exdev()) + _cross_device(monkeypatch) + real_move = shutil.move + copied_mtime_ns = _SOURCE_MTIME_NS + 2_000_000_000 + + def move_onto_coarser_clock(src, dst): + real_move(src, dst) + os.utime(dst, ns=(copied_mtime_ns, copied_mtime_ns)) + + monkeypatch.setattr(ingest_module.shutil, "move", move_onto_coarser_clock) result = _upload(temp) assert result.created_new is True assert dest.read_bytes() == _CONTENT - staging_src, staging_dst = calls[1] - assert Path(staging_dst) == dest - assert Path(staging_src).parent == input_root - assert Path(staging_src).name.startswith(".") - assert Path(staging_src).name.endswith(".tmp") - if sys.platform != "win32": - assert dest.stat().st_mode & 0o777 == 0o640 _assert_nothing_left(temp, input_root, dest) with mock_create_session() as session: content = session.scalars(select(AssetContent)).one() assert content.path == str(dest) assert content.size_bytes == len(_CONTENT) - assert content.mtime_ns == dest.stat().st_mtime_ns != _SOURCE_MTIME_NS - - -def test_destination_that_rejects_mode_bits_still_accepts_the_upload( - mock_create_session, hashing_on, dirs, monkeypatch -): - temp_root, input_root = dirs - temp = _write_temp(temp_root) - dest = _dest_for(input_root, temp) - - def no_chmod(src, dst): - raise PermissionError(errno.EPERM, "Operation not permitted") - - monkeypatch.setattr(ingest_module.shutil, "copymode", no_chmod) - _fail_first_replace(monkeypatch, _exdev()) - - _upload(temp) - - assert dest.read_bytes() == _CONTENT - _assert_nothing_left(temp, input_root, dest) + assert content.mtime_ns == dest.stat().st_mtime_ns == copied_mtime_ns -def test_copy_failure_leaves_no_destination_and_no_staging_file( +def test_failed_cross_device_copy_leaves_no_truncated_file( mock_create_session, hashing_on, dirs, monkeypatch ): temp_root, input_root = dirs temp = _write_temp(temp_root) + _cross_device(monkeypatch) def partial_copy(src, dst): - dst.write(src.read(5)) + Path(dst).write_bytes(_CONTENT[:5]) raise OSError(errno.ENOSPC, "No space left on device") - monkeypatch.setattr(ingest_module.shutil, "copyfileobj", partial_copy) - _fail_first_replace(monkeypatch, _exdev()) + monkeypatch.setattr(ingest_module.shutil, "move", partial_copy) with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): _upload(temp) @@ -160,123 +131,24 @@ def partial_copy(src, dst): _assert_nothing_left(temp, input_root, None) -def _rewrite_same_size(temp: Path) -> None: - temp.write_bytes(_CONTENT[::-1]) - os.utime(temp, ns=(_SOURCE_MTIME_NS + 2_000_000_000,) * 2) - - -def test_source_changed_after_hashing_is_not_copied( - mock_create_session, hashing_on, dirs, monkeypatch -): - temp_root, input_root = dirs - temp = _write_temp(temp_root) - - def rewrite_then_exdev(src, dst): - _rewrite_same_size(temp) - raise _exdev() - - monkeypatch.setattr(ingest_module.os, "replace", rewrite_then_exdev) - - with pytest.raises(UploadUnstableError, match="changed after hashing"): - _upload(temp) - - _assert_nothing_left(temp, input_root, None) - - -def test_source_changed_during_the_copy_is_not_published( - mock_create_session, hashing_on, dirs, monkeypatch -): - temp_root, input_root = dirs - temp = _write_temp(temp_root) - real_copyfileobj = shutil.copyfileobj - - def copy_then_rewrite(src, dst): - real_copyfileobj(src, dst) - _rewrite_same_size(temp) - - monkeypatch.setattr(ingest_module.shutil, "copyfileobj", copy_then_rewrite) - _fail_first_replace(monkeypatch, _exdev()) - - with pytest.raises(UploadUnstableError, match="changed after hashing"): - _upload(temp) - - _assert_nothing_left(temp, input_root, None) - - -def _deny_staging_writes(monkeypatch: pytest.MonkeyPatch, input_root: Path) -> list: - """An ACL that denies writes in the destination: creating a file there raises - PermissionError while os.access still reports the directory writable. The - first os.replace raises EXDEV and switches os.name to "nt", where - tempfile.mkstemp retries that PermissionError up to TMP_MAX times.""" - real_open = os.open - attempts: list = [] - - def denying_open(path, flags, mode=0o777, *, dir_fd=None): - if os.path.dirname(path) == str(input_root): - attempts.append(path) - if len(attempts) > 50: - raise RuntimeError("staging creation kept retrying") - raise PermissionError(errno.EACCES, "Access is denied", str(path)) - return real_open(path, flags, mode, dir_fd=dir_fd) - - def exdev_as_windows(src, dst): - monkeypatch.setattr(os, "name", "nt") - raise _exdev() - - monkeypatch.setattr(ingest_module.os, "open", denying_open) - monkeypatch.setattr(ingest_module.os, "replace", exdev_as_windows) - return attempts - - -def test_write_denied_destination_fails_fast( - mock_create_session, hashing_on, dirs, monkeypatch -): - temp_root, input_root = dirs - temp = _write_temp(temp_root) - attempts = _deny_staging_writes(monkeypatch, input_root) - - with pytest.raises(RuntimeError, match="failed to move uploaded file into place: .*denied"): - _upload(temp) - - monkeypatch.undo() - assert len(attempts) == 1 - _assert_nothing_left(temp, input_root, None) - - -def test_staging_name_collision_takes_another_name( +def test_failed_cross_device_copy_keeps_a_file_already_at_the_destination( mock_create_session, hashing_on, dirs, monkeypatch ): temp_root, input_root = dirs temp = _write_temp(temp_root) dest = _dest_for(input_root, temp) - taken = iter(["0" * 16, "0" * 16, "1" * 16]) - (input_root / f".{'0' * 16}.upload.tmp").write_bytes(b"someone else's") - monkeypatch.setattr(ingest_module.secrets, "token_hex", lambda _n: next(taken)) - _fail_first_replace(monkeypatch, _exdev()) - - _upload(temp) - - assert dest.read_bytes() == _CONTENT - assert sorted(p.name for p in input_root.iterdir()) == sorted( - [dest.name, f".{'0' * 16}.upload.tmp"] - ) + dest.write_bytes(_CONTENT) + _cross_device(monkeypatch) + def unreadable_source(src, dst): + raise PermissionError(errno.EACCES, "Permission denied", src) -def test_final_rename_failure_removes_the_staging_file( - mock_create_session, hashing_on, dirs, monkeypatch -): - temp_root, input_root = dirs - temp = _write_temp(temp_root) - - def locked(src, dst): - raise PermissionError(errno.EACCES, "The process cannot access the file") - - _fail_first_replace(monkeypatch, _exdev(), then=locked) + monkeypatch.setattr(ingest_module.shutil, "move", unreadable_source) with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): _upload(temp) - _assert_nothing_left(temp, input_root, None) + _assert_nothing_left(temp, input_root, dest) def test_other_move_errors_do_not_fall_back_to_a_copy( @@ -284,14 +156,14 @@ def test_other_move_errors_do_not_fall_back_to_a_copy( ): temp_root, input_root = dirs temp = _write_temp(temp_root) - copies: list = [] - monkeypatch.setattr(ingest_module.shutil, "copyfileobj", lambda *a: copies.append(a)) - _fail_first_replace(monkeypatch, PermissionError(errno.EACCES, "Access is denied")) + moves: list = [] + monkeypatch.setattr(ingest_module.shutil, "move", lambda *a: moves.append(a)) + _cross_device(monkeypatch, PermissionError(errno.EACCES, "Access is denied")) with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): _upload(temp) - assert copies == [] + assert moves == [] _assert_nothing_left(temp, input_root, None) @@ -301,12 +173,12 @@ def test_same_device_upload_is_a_plain_rename( temp_root, input_root = dirs temp = _write_temp(temp_root) dest = _dest_for(input_root, temp) - copies: list = [] - monkeypatch.setattr(ingest_module.shutil, "copyfileobj", lambda *a: copies.append(a)) + moves: list = [] + monkeypatch.setattr(ingest_module.shutil, "move", lambda *a: moves.append(a)) _upload(temp) - assert copies == [] + assert moves == [] assert dest.read_bytes() == _CONTENT _assert_nothing_left(temp, input_root, dest) with mock_create_session() as session: From 29471cc9a9877be52bb593267e34ecb2225f5720 Mon Sep 17 00:00:00 2001 From: Simon Pinfold Date: Fri, 2 Oct 2026 20:20:11 -0700 Subject: [PATCH 07/15] fix(assets): stage cross-volume copies beside the destination Copying straight onto the hash-named path overwrote an existing file in place, so a copy that failed partway (e.g. a full disk) left a good asset truncated. shutil.move into a uniquely named .part sibling and os.replace it into place, removing the .part on failure. --- app/assets/services/ingest.py | 18 ++++++++++-------- .../services/test_upload_cross_device.py | 19 +++++++++++++------ 2 files changed, 23 insertions(+), 14 deletions(-) diff --git a/app/assets/services/ingest.py b/app/assets/services/ingest.py index 20bbb56eac4..abcbe06ab8f 100644 --- a/app/assets/services/ingest.py +++ b/app/assets/services/ingest.py @@ -14,6 +14,7 @@ import mimetypes import os import shutil +import uuid from typing import Any, NamedTuple, Sequence from sqlalchemy import func, select @@ -192,7 +193,7 @@ def _move_temp_to_dest( ) -> os.stat_result: """Move the upload into place and return the stat to record for it. Across volumes (EXDEV, also Windows' ERROR_NOT_SAME_DEVICE) the move is a copy with - its own mtime, so that case records the destination's stat.""" + its own mtime, so that case records the copy's stat.""" os.makedirs(os.path.dirname(dest_abs), exist_ok=True) try: try: @@ -201,16 +202,17 @@ def _move_temp_to_dest( except OSError as e: if e.errno != errno.EXDEV: raise - dest_existed = os.path.exists(dest_abs) + # Copy beside dest and rename, so a partial copy never takes the hash name. + staging = f"{dest_abs}.{uuid.uuid4().hex}.part" try: - shutil.move(temp_path, dest_abs) + shutil.move(temp_path, staging) + staged_stat = os.stat(staging) + os.replace(staging, dest_abs) except BaseException: - # A failed copy (e.g. a full disk) leaves a truncated file at dest. - if not dest_existed: - with contextlib.suppress(OSError): - os.remove(dest_abs) + with contextlib.suppress(OSError): + os.remove(staging) raise - return os.stat(dest_abs) + return staged_stat except Exception as e: raise RuntimeError(f"failed to move uploaded file into place: {e}") from e diff --git a/tests-unit/assets_test/services/test_upload_cross_device.py b/tests-unit/assets_test/services/test_upload_cross_device.py index 52deb62839a..4c599e0038c 100644 --- a/tests-unit/assets_test/services/test_upload_cross_device.py +++ b/tests-unit/assets_test/services/test_upload_cross_device.py @@ -66,11 +66,15 @@ def _upload(temp: Path): def _cross_device(monkeypatch: pytest.MonkeyPatch, exc: OSError | None = None) -> None: - """Make renames fail as they do across volumes: os.replace raises ``exc`` - (EXDEV by default), and so does the os.rename shutil.move tries first.""" + """Make renames between directories fail as they do across volumes: os.replace + raises ``exc`` (EXDEV by default), and so does the os.rename shutil.move tries + first. Renames within one directory (staging file to final name) still work.""" exc = exc or OSError(errno.EXDEV, "Invalid cross-device link") + real_replace = os.replace def fail(src, dst): + if os.path.dirname(src) == os.path.dirname(dst): + return real_replace(src, dst) raise exc monkeypatch.setattr(ingest_module.os, "replace", fail) @@ -131,7 +135,7 @@ def partial_copy(src, dst): _assert_nothing_left(temp, input_root, None) -def test_failed_cross_device_copy_keeps_a_file_already_at_the_destination( +def test_failed_cross_device_copy_leaves_an_existing_file_untouched( mock_create_session, hashing_on, dirs, monkeypatch ): temp_root, input_root = dirs @@ -140,14 +144,17 @@ def test_failed_cross_device_copy_keeps_a_file_already_at_the_destination( dest.write_bytes(_CONTENT) _cross_device(monkeypatch) - def unreadable_source(src, dst): - raise PermissionError(errno.EACCES, "Permission denied", src) + def disk_full_mid_copy(src, dst): + with open(dst, "wb") as partial: + partial.write(_CONTENT[:5]) + raise OSError(errno.ENOSPC, "No space left on device") - monkeypatch.setattr(ingest_module.shutil, "move", unreadable_source) + monkeypatch.setattr(ingest_module.shutil, "move", disk_full_mid_copy) with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): _upload(temp) + assert dest.read_bytes() == _CONTENT _assert_nothing_left(temp, input_root, dest) From e6563d57b5dd3607bf0a4c9f338b5dd9e2213a50 Mon Sep 17 00:00:00 2001 From: Simon Pinfold Date: Fri, 2 Oct 2026 20:24:05 -0700 Subject: [PATCH 08/15] refactor(assets): reduce the cross-volume upload fallback to a plain copy On EXDEV, copy the upload into place with shutil.copyfile, skipping the copy when the hash-named file is already there at the same size, and record the stat of the file at the destination. Drop the staging file and its tests. --- app/assets/services/ingest.py | 36 +-- .../services/test_upload_cross_device.py | 211 +++--------------- 2 files changed, 38 insertions(+), 209 deletions(-) diff --git a/app/assets/services/ingest.py b/app/assets/services/ingest.py index abcbe06ab8f..7649ef60087 100644 --- a/app/assets/services/ingest.py +++ b/app/assets/services/ingest.py @@ -1,11 +1,10 @@ """Turns incoming bytes into catalogued assets: multipart uploads moved into a hash-addressed destination, files registered where they already sit, and records created from a hash the catalog already holds. Every path persists the -stat that hashing verified (for an upload copied across volumes, the stat of the -copy), so a row's recorded size and mtime describe the same observation as its -hash. A live row already at the destination is reconciled before the write, so -an upload never adopts a fresh hash onto records created for bytes it just -replaced. +stat that hashing verified, so a row's recorded size and mtime describe the +same observation as its hash. A live row already at the destination is +reconciled before the write, so an upload never adopts a fresh hash onto +records created for bytes it just replaced. """ import contextlib @@ -14,7 +13,6 @@ import mimetypes import os import shutil -import uuid from typing import Any, NamedTuple, Sequence from sqlalchemy import func, select @@ -188,31 +186,16 @@ def _guess_upload_mime_type( return guessed or "application/octet-stream" -def _move_temp_to_dest( - temp_path: str, dest_abs: str, verified_stat: os.stat_result -) -> os.stat_result: - """Move the upload into place and return the stat to record for it. Across - volumes (EXDEV, also Windows' ERROR_NOT_SAME_DEVICE) the move is a copy with - its own mtime, so that case records the copy's stat.""" +def _move_temp_to_dest(temp_path: str, dest_abs: str) -> None: os.makedirs(os.path.dirname(dest_abs), exist_ok=True) try: try: os.replace(temp_path, dest_abs) - return verified_stat - except OSError as e: + except OSError as e: # EXDEV: destination is on another volume if e.errno != errno.EXDEV: raise - # Copy beside dest and rename, so a partial copy never takes the hash name. - staging = f"{dest_abs}.{uuid.uuid4().hex}.part" - try: - shutil.move(temp_path, staging) - staged_stat = os.stat(staging) - os.replace(staging, dest_abs) - except BaseException: - with contextlib.suppress(OSError): - os.remove(staging) - raise - return staged_stat + if not (os.path.exists(dest_abs) and os.path.getsize(dest_abs) == os.path.getsize(temp_path)): + shutil.copyfile(temp_path, dest_abs) except Exception as e: raise RuntimeError(f"failed to move uploaded file into place: {e}") from e @@ -511,9 +494,10 @@ def upload_from_temp_path( content_type = _guess_upload_mime_type( mime_type, client_filename, name, os.path.basename(dest_abs) ) - placed_stat = _move_temp_to_dest(temp_path, dest_abs, verified_stat) + _move_temp_to_dest(temp_path, dest_abs) finally: _remove_temp_path(temp_path) + placed_stat = os.stat(dest_abs) size_bytes, mtime_ns = placed_stat.st_size, placed_stat.st_mtime_ns system_metadata = _extract_system_metadata_sync(dest_abs, content_type) with create_session() as session: diff --git a/tests-unit/assets_test/services/test_upload_cross_device.py b/tests-unit/assets_test/services/test_upload_cross_device.py index 4c599e0038c..72de5e645d7 100644 --- a/tests-unit/assets_test/services/test_upload_cross_device.py +++ b/tests-unit/assets_test/services/test_upload_cross_device.py @@ -1,226 +1,71 @@ -"""Uploads whose destination sits on a different volume than the temp upload -(another drive letter on Windows, another mount on POSIX), where a rename fails -with EXDEV and the bytes have to be copied instead.""" +"""Uploads whose destination is on a different volume than the temp upload, where +the rename fails with EXDEV (WinError 17 on Windows) and the bytes are copied.""" import errno import os -import shutil -import sys import uuid from pathlib import Path import pytest -from sqlalchemy import select import app.assets.mode as mode_module import app.assets.services.ingest as ingest_module import folder_paths -from app.assets.database.models import AssetContent from app.assets.services.ingest import upload_from_temp_path from app.assets.services.snapshot_hash import snapshot_hash _CONTENT = b"cross-device upload bytes" -_SOURCE_MTIME_NS = 1_600_000_000_000_000_000 # whole seconds: NTFS keeps 100 ns @pytest.fixture -def hashing_on(): +def cross_device_upload(tmp_path: Path, monkeypatch: pytest.MonkeyPatch): class FakeArgs: enable_asset_hashing = True mode_module.init(FakeArgs()) - yield + monkeypatch.setattr(folder_paths, "get_temp_directory", lambda: str(tmp_path / "temp")) + monkeypatch.setattr(folder_paths, "get_input_directory", lambda: str(tmp_path / "input")) + + def exdev(src, dst): + raise OSError(errno.EXDEV, "Invalid cross-device link") + + monkeypatch.setattr(ingest_module.os, "replace", exdev) + temp = tmp_path / "temp" / "uploads" / uuid.uuid4().hex / ".upload.part" + temp.parent.mkdir(parents=True) + temp.write_bytes(_CONTENT) + dest = tmp_path / "input" / f"{snapshot_hash(str(temp))[0]}.png" + yield temp, dest mode_module.init(None) -@pytest.fixture -def dirs(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> tuple[Path, Path]: - temp_root = tmp_path / "temp" - input_root = tmp_path / "input" - temp_root.mkdir() - input_root.mkdir() - monkeypatch.setattr(folder_paths, "get_temp_directory", lambda: str(temp_root)) - monkeypatch.setattr(folder_paths, "get_input_directory", lambda: str(input_root)) - return temp_root, input_root - - -def _write_temp(temp_root: Path) -> Path: - unique_dir = temp_root / "uploads" / uuid.uuid4().hex - unique_dir.mkdir(parents=True) - path = unique_dir / ".upload.part" - path.write_bytes(_CONTENT) - os.utime(path, ns=(_SOURCE_MTIME_NS, _SOURCE_MTIME_NS)) - return path - - -def _dest_for(input_root: Path, temp: Path) -> Path: - snapshot = snapshot_hash(str(temp)) - assert snapshot is not None - return input_root / f"{snapshot[0]}.png" - - def _upload(temp: Path): return upload_from_temp_path( temp_path=str(temp), name="photo.png", tags=["input"], client_filename="photo.png" ) -def _cross_device(monkeypatch: pytest.MonkeyPatch, exc: OSError | None = None) -> None: - """Make renames between directories fail as they do across volumes: os.replace - raises ``exc`` (EXDEV by default), and so does the os.rename shutil.move tries - first. Renames within one directory (staging file to final name) still work.""" - exc = exc or OSError(errno.EXDEV, "Invalid cross-device link") - real_replace = os.replace - - def fail(src, dst): - if os.path.dirname(src) == os.path.dirname(dst): - return real_replace(src, dst) - raise exc - - monkeypatch.setattr(ingest_module.os, "replace", fail) - monkeypatch.setattr(ingest_module.os, "rename", fail) - - -def _assert_nothing_left(temp: Path, input_root: Path, dest: Path | None) -> None: - assert not temp.exists() - assert not temp.parent.exists() - leftovers = sorted(p.name for p in input_root.iterdir()) - assert leftovers == ([dest.name] if dest is not None else []) - - -def test_cross_device_upload_is_copied_into_place( - mock_create_session, hashing_on, dirs, monkeypatch -): - temp_root, input_root = dirs - temp = _write_temp(temp_root) - dest = _dest_for(input_root, temp) - _cross_device(monkeypatch) - real_move = shutil.move - copied_mtime_ns = _SOURCE_MTIME_NS + 2_000_000_000 - - def move_onto_coarser_clock(src, dst): - real_move(src, dst) - os.utime(dst, ns=(copied_mtime_ns, copied_mtime_ns)) - - monkeypatch.setattr(ingest_module.shutil, "move", move_onto_coarser_clock) +def test_cross_device_upload_is_copied_into_place(mock_create_session, cross_device_upload): + temp, dest = cross_device_upload result = _upload(temp) assert result.created_new is True assert dest.read_bytes() == _CONTENT - _assert_nothing_left(temp, input_root, dest) - with mock_create_session() as session: - content = session.scalars(select(AssetContent)).one() - assert content.path == str(dest) - assert content.size_bytes == len(_CONTENT) - assert content.mtime_ns == dest.stat().st_mtime_ns == copied_mtime_ns - - -def test_failed_cross_device_copy_leaves_no_truncated_file( - mock_create_session, hashing_on, dirs, monkeypatch -): - temp_root, input_root = dirs - temp = _write_temp(temp_root) - _cross_device(monkeypatch) - - def partial_copy(src, dst): - Path(dst).write_bytes(_CONTENT[:5]) - raise OSError(errno.ENOSPC, "No space left on device") - - monkeypatch.setattr(ingest_module.shutil, "move", partial_copy) - - with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): - _upload(temp) - - _assert_nothing_left(temp, input_root, None) + assert not temp.exists() -def test_failed_cross_device_copy_leaves_an_existing_file_untouched( - mock_create_session, hashing_on, dirs, monkeypatch +def test_cross_device_upload_keeps_a_same_size_file_already_there( + mock_create_session, cross_device_upload, monkeypatch ): - temp_root, input_root = dirs - temp = _write_temp(temp_root) - dest = _dest_for(input_root, temp) + temp, dest = cross_device_upload + dest.parent.mkdir(parents=True) dest.write_bytes(_CONTENT) - _cross_device(monkeypatch) - - def disk_full_mid_copy(src, dst): - with open(dst, "wb") as partial: - partial.write(_CONTENT[:5]) - raise OSError(errno.ENOSPC, "No space left on device") - - monkeypatch.setattr(ingest_module.shutil, "move", disk_full_mid_copy) - - with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): - _upload(temp) - - assert dest.read_bytes() == _CONTENT - _assert_nothing_left(temp, input_root, dest) - - -def test_other_move_errors_do_not_fall_back_to_a_copy( - mock_create_session, hashing_on, dirs, monkeypatch -): - temp_root, input_root = dirs - temp = _write_temp(temp_root) - moves: list = [] - monkeypatch.setattr(ingest_module.shutil, "move", lambda *a: moves.append(a)) - _cross_device(monkeypatch, PermissionError(errno.EACCES, "Access is denied")) - - with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): - _upload(temp) - - assert moves == [] - _assert_nothing_left(temp, input_root, None) - - -def test_same_device_upload_is_a_plain_rename( - mock_create_session, hashing_on, dirs, monkeypatch -): - temp_root, input_root = dirs - temp = _write_temp(temp_root) - dest = _dest_for(input_root, temp) - moves: list = [] - monkeypatch.setattr(ingest_module.shutil, "move", lambda *a: moves.append(a)) + os.utime(dest, ns=(1_600_000_000_000_000_000,) * 2) + copies: list = [] + monkeypatch.setattr(ingest_module.shutil, "copyfile", lambda *a: copies.append(a)) _upload(temp) - assert moves == [] - assert dest.read_bytes() == _CONTENT - _assert_nothing_left(temp, input_root, dest) - with mock_create_session() as session: - assert session.scalars(select(AssetContent)).one().mtime_ns == _SOURCE_MTIME_NS - - -@pytest.mark.skipif(sys.platform != "win32", reason="Windows error mapping") -def test_windows_not_same_device_error_maps_to_exdev(): - # ERROR_NOT_SAME_DEVICE (17) is what os.replace raises across drive letters. - assert OSError(None, "not same device", None, 17).errno == errno.EXDEV - - -def test_real_cross_filesystem_upload(mock_create_session, hashing_on, tmp_path, monkeypatch): - shm = Path("/dev/shm") - if not shm.is_dir() or not os.access(shm, os.W_OK): - pytest.skip("no writable /dev/shm") - if shm.stat().st_dev == tmp_path.stat().st_dev: - pytest.skip("/dev/shm and tmp_path share a filesystem") - temp_root = shm / f"comfy-cross-device-{uuid.uuid4().hex}" - input_root = tmp_path / "input" - temp_root.mkdir() - input_root.mkdir() - monkeypatch.setattr(folder_paths, "get_temp_directory", lambda: str(temp_root)) - monkeypatch.setattr(folder_paths, "get_input_directory", lambda: str(input_root)) - try: - temp = _write_temp(temp_root) - dest = _dest_for(input_root, temp) - - _upload(temp) - - assert dest.read_bytes() == _CONTENT - _assert_nothing_left(temp, input_root, dest) - with mock_create_session() as session: - assert session.scalars(select(AssetContent)).one().mtime_ns == ( - dest.stat().st_mtime_ns - ) - finally: - shutil.rmtree(temp_root, ignore_errors=True) + assert copies == [] + assert dest.stat().st_mtime_ns == 1_600_000_000_000_000_000 + assert not temp.exists() From 75f56d81491278f6606c5af95613f950d06d6950 Mon Sep 17 00:00:00 2001 From: Simon Pinfold Date: Fri, 2 Oct 2026 20:26:35 -0700 Subject: [PATCH 09/15] docs(assets): note why the placed file's stat is recorded --- app/assets/services/ingest.py | 1 + 1 file changed, 1 insertion(+) diff --git a/app/assets/services/ingest.py b/app/assets/services/ingest.py index 7649ef60087..ec04c8d8b6a 100644 --- a/app/assets/services/ingest.py +++ b/app/assets/services/ingest.py @@ -497,6 +497,7 @@ def upload_from_temp_path( _move_temp_to_dest(temp_path, dest_abs) finally: _remove_temp_path(temp_path) + # A cross-volume copy gets a new mtime, so record the file on disk (a rename keeps it). placed_stat = os.stat(dest_abs) size_bytes, mtime_ns = placed_stat.st_size, placed_stat.st_mtime_ns system_metadata = _extract_system_metadata_sync(dest_abs, content_type) From 8160762001d3804915d50635fafe72a917c0dbb7 Mon Sep 17 00:00:00 2001 From: Simon Pinfold Date: Sat, 3 Oct 2026 23:53:15 -0700 Subject: [PATCH 10/15] fix(assets): compare content, not size, before keeping a cross-volume destination Keep an existing hash-named file only when its own hash matches the upload; otherwise remove the entry first, so the copy never writes through a link or over other bytes. Tests now pin the recorded mtime, the replacement of mismatched or truncated files, links, and non-EXDEV errors. --- app/assets/services/ingest.py | 20 +++--- .../services/test_upload_cross_device.py | 72 ++++++++++++++++++- 2 files changed, 81 insertions(+), 11 deletions(-) diff --git a/app/assets/services/ingest.py b/app/assets/services/ingest.py index ec04c8d8b6a..b48163380a6 100644 --- a/app/assets/services/ingest.py +++ b/app/assets/services/ingest.py @@ -1,10 +1,11 @@ """Turns incoming bytes into catalogued assets: multipart uploads moved into a hash-addressed destination, files registered where they already sit, and -records created from a hash the catalog already holds. Every path persists the -stat that hashing verified, so a row's recorded size and mtime describe the -same observation as its hash. A live row already at the destination is -reconciled before the write, so an upload never adopts a fresh hash onto -records created for bytes it just replaced. +records created from a hash the catalog already holds. Registration persists +the stat that hashing verified, so a row's recorded size and mtime describe the +same observation as its hash; an upload records the stat of the file once it is +in place, since a copy across volumes has its own mtime. A live row already at +the destination is reconciled before the write, so an upload never adopts a +fresh hash onto records created for bytes it just replaced. """ import contextlib @@ -186,7 +187,7 @@ def _guess_upload_mime_type( return guessed or "application/octet-stream" -def _move_temp_to_dest(temp_path: str, dest_abs: str) -> None: +def _move_temp_to_dest(temp_path: str, dest_abs: str, digest: str) -> None: os.makedirs(os.path.dirname(dest_abs), exist_ok=True) try: try: @@ -194,7 +195,10 @@ def _move_temp_to_dest(temp_path: str, dest_abs: str) -> None: except OSError as e: # EXDEV: destination is on another volume if e.errno != errno.EXDEV: raise - if not (os.path.exists(dest_abs) and os.path.getsize(dest_abs) == os.path.getsize(temp_path)): + existing = snapshot_hash(dest_abs) + if existing is None or existing[0] != digest: + with contextlib.suppress(FileNotFoundError): + os.remove(dest_abs) # never write through a link or over other bytes shutil.copyfile(temp_path, dest_abs) except Exception as e: raise RuntimeError(f"failed to move uploaded file into place: {e}") from e @@ -494,7 +498,7 @@ def upload_from_temp_path( content_type = _guess_upload_mime_type( mime_type, client_filename, name, os.path.basename(dest_abs) ) - _move_temp_to_dest(temp_path, dest_abs) + _move_temp_to_dest(temp_path, dest_abs, digest) finally: _remove_temp_path(temp_path) # A cross-volume copy gets a new mtime, so record the file on disk (a rename keeps it). diff --git a/tests-unit/assets_test/services/test_upload_cross_device.py b/tests-unit/assets_test/services/test_upload_cross_device.py index 72de5e645d7..a33ac071479 100644 --- a/tests-unit/assets_test/services/test_upload_cross_device.py +++ b/tests-unit/assets_test/services/test_upload_cross_device.py @@ -7,14 +7,17 @@ from pathlib import Path import pytest +from sqlalchemy import select import app.assets.mode as mode_module import app.assets.services.ingest as ingest_module import folder_paths +from app.assets.database.models import AssetContent from app.assets.services.ingest import upload_from_temp_path from app.assets.services.snapshot_hash import snapshot_hash _CONTENT = b"cross-device upload bytes" +_OLD_MTIME_NS = 1_600_000_000_000_000_000 @pytest.fixture @@ -33,6 +36,7 @@ def exdev(src, dst): temp = tmp_path / "temp" / "uploads" / uuid.uuid4().hex / ".upload.part" temp.parent.mkdir(parents=True) temp.write_bytes(_CONTENT) + os.utime(temp, ns=(_OLD_MTIME_NS, _OLD_MTIME_NS)) dest = tmp_path / "input" / f"{snapshot_hash(str(temp))[0]}.png" yield temp, dest mode_module.init(None) @@ -44,6 +48,13 @@ def _upload(temp: Path): ) +def _recorded_mtime_ns(session_factory, dest: Path) -> int: + with session_factory() as session: + return session.scalars( + select(AssetContent.mtime_ns).where(AssetContent.path == str(dest)) + ).one() + + def test_cross_device_upload_is_copied_into_place(mock_create_session, cross_device_upload): temp, dest = cross_device_upload @@ -52,20 +63,75 @@ def test_cross_device_upload_is_copied_into_place(mock_create_session, cross_dev assert result.created_new is True assert dest.read_bytes() == _CONTENT assert not temp.exists() + # The copy has its own mtime, and that is what must be recorded. + assert dest.stat().st_mtime_ns != _OLD_MTIME_NS + assert _recorded_mtime_ns(mock_create_session, dest) == dest.stat().st_mtime_ns -def test_cross_device_upload_keeps_a_same_size_file_already_there( +def test_cross_device_upload_keeps_identical_bytes_already_there( mock_create_session, cross_device_upload, monkeypatch ): temp, dest = cross_device_upload dest.parent.mkdir(parents=True) dest.write_bytes(_CONTENT) - os.utime(dest, ns=(1_600_000_000_000_000_000,) * 2) + existing_mtime_ns = _OLD_MTIME_NS - 1_000_000_000 + os.utime(dest, ns=(existing_mtime_ns, existing_mtime_ns)) copies: list = [] monkeypatch.setattr(ingest_module.shutil, "copyfile", lambda *a: copies.append(a)) _upload(temp) assert copies == [] - assert dest.stat().st_mtime_ns == 1_600_000_000_000_000_000 + assert dest.stat().st_mtime_ns == existing_mtime_ns + assert _recorded_mtime_ns(mock_create_session, dest) == existing_mtime_ns assert not temp.exists() + + +@pytest.mark.parametrize( + "existing", [_CONTENT[::-1], _CONTENT[:5]], ids=["same-size-other-bytes", "truncated"] +) +def test_cross_device_upload_replaces_other_bytes_at_the_destination( + mock_create_session, cross_device_upload, existing +): + temp, dest = cross_device_upload + dest.parent.mkdir(parents=True) + dest.write_bytes(existing) + + _upload(temp) + + assert dest.read_bytes() == _CONTENT + + +@pytest.mark.skipif(os.name == "nt", reason="symlinks need privileges on Windows") +def test_cross_device_upload_does_not_write_through_a_link( + mock_create_session, cross_device_upload, tmp_path +): + temp, dest = cross_device_upload + dest.parent.mkdir(parents=True) + elsewhere = tmp_path / "elsewhere.png" + elsewhere.write_bytes(b"someone else's file") + dest.symlink_to(elsewhere) + + _upload(temp) + + assert elsewhere.read_bytes() == b"someone else's file" + assert not dest.is_symlink() + assert dest.read_bytes() == _CONTENT + + +def test_other_move_errors_raise_without_copying( + mock_create_session, cross_device_upload, monkeypatch +): + temp, _dest = cross_device_upload + + def denied(src, dst): + raise PermissionError(errno.EACCES, "Access is denied") + + monkeypatch.setattr(ingest_module.os, "replace", denied) + copies: list = [] + monkeypatch.setattr(ingest_module.shutil, "copyfile", lambda *a: copies.append(a)) + + with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): + _upload(temp) + + assert copies == [] From 67675c0275a36ae50fac614bdd24ab88d9b1690b Mon Sep 17 00:00:00 2001 From: Simon Pinfold Date: Sun, 4 Oct 2026 00:03:12 -0700 Subject: [PATCH 11/15] test(assets): cover a dangling symlink at a cross-volume destination --- .../assets_test/services/test_upload_cross_device.py | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/tests-unit/assets_test/services/test_upload_cross_device.py b/tests-unit/assets_test/services/test_upload_cross_device.py index a33ac071479..64346d56e32 100644 --- a/tests-unit/assets_test/services/test_upload_cross_device.py +++ b/tests-unit/assets_test/services/test_upload_cross_device.py @@ -103,18 +103,23 @@ def test_cross_device_upload_replaces_other_bytes_at_the_destination( @pytest.mark.skipif(os.name == "nt", reason="symlinks need privileges on Windows") +@pytest.mark.parametrize("target_exists", [True, False], ids=["other-file", "dangling"]) def test_cross_device_upload_does_not_write_through_a_link( - mock_create_session, cross_device_upload, tmp_path + mock_create_session, cross_device_upload, tmp_path, target_exists ): temp, dest = cross_device_upload dest.parent.mkdir(parents=True) elsewhere = tmp_path / "elsewhere.png" - elsewhere.write_bytes(b"someone else's file") + if target_exists: + elsewhere.write_bytes(b"someone else's file") dest.symlink_to(elsewhere) _upload(temp) - assert elsewhere.read_bytes() == b"someone else's file" + if target_exists: + assert elsewhere.read_bytes() == b"someone else's file" + else: + assert not elsewhere.exists() assert not dest.is_symlink() assert dest.read_bytes() == _CONTENT From 1f3b60d657e4f9d6b0ceb0d7a26de7dfc9168ba0 Mon Sep 17 00:00:00 2001 From: Simon Pinfold Date: Sun, 4 Oct 2026 00:29:56 -0700 Subject: [PATCH 12/15] fix(assets): replace linked destinations and clean up a failed cross-volume copy Keep an existing hash-named file only when it is a regular, single-link file with matching bytes, so a symlink or hardlink left by a dedupe tool is replaced as a same-volume rename would. Remove the partial file if the copy fails, so a full disk doesn't stay full. --- app/assets/services/ingest.py | 12 +++-- .../services/test_upload_cross_device.py | 53 ++++++++++++++++--- 2 files changed, 54 insertions(+), 11 deletions(-) diff --git a/app/assets/services/ingest.py b/app/assets/services/ingest.py index b48163380a6..05c1cec6718 100644 --- a/app/assets/services/ingest.py +++ b/app/assets/services/ingest.py @@ -195,11 +195,17 @@ def _move_temp_to_dest(temp_path: str, dest_abs: str, digest: str) -> None: except OSError as e: # EXDEV: destination is on another volume if e.errno != errno.EXDEV: raise - existing = snapshot_hash(dest_abs) - if existing is None or existing[0] != digest: + # Keep only a regular, single-link file with these exact bytes; replace anything else. + existing = None if os.path.islink(dest_abs) else snapshot_hash(dest_abs) + if existing is None or existing[0] != digest or existing[1].st_nlink != 1: with contextlib.suppress(FileNotFoundError): os.remove(dest_abs) # never write through a link or over other bytes - shutil.copyfile(temp_path, dest_abs) + try: + shutil.copyfile(temp_path, dest_abs) + except BaseException: + with contextlib.suppress(OSError): + os.remove(dest_abs) # don't leave a partial copy under the hash name + raise except Exception as e: raise RuntimeError(f"failed to move uploaded file into place: {e}") from e diff --git a/tests-unit/assets_test/services/test_upload_cross_device.py b/tests-unit/assets_test/services/test_upload_cross_device.py index 64346d56e32..0939918ee83 100644 --- a/tests-unit/assets_test/services/test_upload_cross_device.py +++ b/tests-unit/assets_test/services/test_upload_cross_device.py @@ -103,27 +103,64 @@ def test_cross_device_upload_replaces_other_bytes_at_the_destination( @pytest.mark.skipif(os.name == "nt", reason="symlinks need privileges on Windows") -@pytest.mark.parametrize("target_exists", [True, False], ids=["other-file", "dangling"]) -def test_cross_device_upload_does_not_write_through_a_link( - mock_create_session, cross_device_upload, tmp_path, target_exists +@pytest.mark.parametrize( + "target", [b"someone else's file", _CONTENT, None], ids=["other-bytes", "same-bytes", "dangling"] +) +def test_cross_device_upload_replaces_a_symlink( + mock_create_session, cross_device_upload, tmp_path, target ): temp, dest = cross_device_upload dest.parent.mkdir(parents=True) elsewhere = tmp_path / "elsewhere.png" - if target_exists: - elsewhere.write_bytes(b"someone else's file") + if target is not None: + elsewhere.write_bytes(target) dest.symlink_to(elsewhere) _upload(temp) - if target_exists: - assert elsewhere.read_bytes() == b"someone else's file" - else: + if target is None: assert not elsewhere.exists() + else: + assert elsewhere.read_bytes() == target assert not dest.is_symlink() assert dest.read_bytes() == _CONTENT +@pytest.mark.parametrize("twin", [b"someone else's file!!", _CONTENT], ids=["other-bytes", "same-bytes"]) +def test_cross_device_upload_replaces_a_hardlink( + mock_create_session, cross_device_upload, tmp_path, twin +): + temp, dest = cross_device_upload + dest.parent.mkdir(parents=True) + elsewhere = tmp_path / "elsewhere.png" + elsewhere.write_bytes(twin) + os.link(elsewhere, dest) + + _upload(temp) + + assert elsewhere.read_bytes() == twin + assert not os.path.samefile(elsewhere, dest) + assert dest.read_bytes() == _CONTENT + + +def test_failed_cross_device_copy_leaves_no_partial_file( + mock_create_session, cross_device_upload, monkeypatch +): + temp, dest = cross_device_upload + + def disk_full(src, dst): + Path(dst).write_bytes(_CONTENT[:5]) + raise OSError(errno.ENOSPC, "No space left on device") + + monkeypatch.setattr(ingest_module.shutil, "copyfile", disk_full) + + with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): + _upload(temp) + + assert not dest.exists() + assert not temp.exists() + + def test_other_move_errors_raise_without_copying( mock_create_session, cross_device_upload, monkeypatch ): From 1bf4d33c0a37abe5cc1e79578b3bd995c007abc5 Mon Sep 17 00:00:00 2001 From: Simon Pinfold Date: Sun, 4 Oct 2026 00:44:32 -0700 Subject: [PATCH 13/15] refactor(assets): reduce the cross-volume fallback to a bare copy On EXDEV, copy the upload onto the hash-named path with shutil.copyfile and record the placed file's stat. Drop the content-hash keep check, the link checks and the partial-file cleanup; their cases are listed as known limitations on the PR. --- app/assets/services/ingest.py | 16 +--- .../services/test_upload_cross_device.py | 78 ------------------- 2 files changed, 3 insertions(+), 91 deletions(-) diff --git a/app/assets/services/ingest.py b/app/assets/services/ingest.py index 05c1cec6718..dccbf02c188 100644 --- a/app/assets/services/ingest.py +++ b/app/assets/services/ingest.py @@ -187,7 +187,7 @@ def _guess_upload_mime_type( return guessed or "application/octet-stream" -def _move_temp_to_dest(temp_path: str, dest_abs: str, digest: str) -> None: +def _move_temp_to_dest(temp_path: str, dest_abs: str) -> None: os.makedirs(os.path.dirname(dest_abs), exist_ok=True) try: try: @@ -195,17 +195,7 @@ def _move_temp_to_dest(temp_path: str, dest_abs: str, digest: str) -> None: except OSError as e: # EXDEV: destination is on another volume if e.errno != errno.EXDEV: raise - # Keep only a regular, single-link file with these exact bytes; replace anything else. - existing = None if os.path.islink(dest_abs) else snapshot_hash(dest_abs) - if existing is None or existing[0] != digest or existing[1].st_nlink != 1: - with contextlib.suppress(FileNotFoundError): - os.remove(dest_abs) # never write through a link or over other bytes - try: - shutil.copyfile(temp_path, dest_abs) - except BaseException: - with contextlib.suppress(OSError): - os.remove(dest_abs) # don't leave a partial copy under the hash name - raise + shutil.copyfile(temp_path, dest_abs) except Exception as e: raise RuntimeError(f"failed to move uploaded file into place: {e}") from e @@ -504,7 +494,7 @@ def upload_from_temp_path( content_type = _guess_upload_mime_type( mime_type, client_filename, name, os.path.basename(dest_abs) ) - _move_temp_to_dest(temp_path, dest_abs, digest) + _move_temp_to_dest(temp_path, dest_abs) finally: _remove_temp_path(temp_path) # A cross-volume copy gets a new mtime, so record the file on disk (a rename keeps it). diff --git a/tests-unit/assets_test/services/test_upload_cross_device.py b/tests-unit/assets_test/services/test_upload_cross_device.py index 0939918ee83..c869c22a8e2 100644 --- a/tests-unit/assets_test/services/test_upload_cross_device.py +++ b/tests-unit/assets_test/services/test_upload_cross_device.py @@ -68,25 +68,6 @@ def test_cross_device_upload_is_copied_into_place(mock_create_session, cross_dev assert _recorded_mtime_ns(mock_create_session, dest) == dest.stat().st_mtime_ns -def test_cross_device_upload_keeps_identical_bytes_already_there( - mock_create_session, cross_device_upload, monkeypatch -): - temp, dest = cross_device_upload - dest.parent.mkdir(parents=True) - dest.write_bytes(_CONTENT) - existing_mtime_ns = _OLD_MTIME_NS - 1_000_000_000 - os.utime(dest, ns=(existing_mtime_ns, existing_mtime_ns)) - copies: list = [] - monkeypatch.setattr(ingest_module.shutil, "copyfile", lambda *a: copies.append(a)) - - _upload(temp) - - assert copies == [] - assert dest.stat().st_mtime_ns == existing_mtime_ns - assert _recorded_mtime_ns(mock_create_session, dest) == existing_mtime_ns - assert not temp.exists() - - @pytest.mark.parametrize( "existing", [_CONTENT[::-1], _CONTENT[:5]], ids=["same-size-other-bytes", "truncated"] ) @@ -102,65 +83,6 @@ def test_cross_device_upload_replaces_other_bytes_at_the_destination( assert dest.read_bytes() == _CONTENT -@pytest.mark.skipif(os.name == "nt", reason="symlinks need privileges on Windows") -@pytest.mark.parametrize( - "target", [b"someone else's file", _CONTENT, None], ids=["other-bytes", "same-bytes", "dangling"] -) -def test_cross_device_upload_replaces_a_symlink( - mock_create_session, cross_device_upload, tmp_path, target -): - temp, dest = cross_device_upload - dest.parent.mkdir(parents=True) - elsewhere = tmp_path / "elsewhere.png" - if target is not None: - elsewhere.write_bytes(target) - dest.symlink_to(elsewhere) - - _upload(temp) - - if target is None: - assert not elsewhere.exists() - else: - assert elsewhere.read_bytes() == target - assert not dest.is_symlink() - assert dest.read_bytes() == _CONTENT - - -@pytest.mark.parametrize("twin", [b"someone else's file!!", _CONTENT], ids=["other-bytes", "same-bytes"]) -def test_cross_device_upload_replaces_a_hardlink( - mock_create_session, cross_device_upload, tmp_path, twin -): - temp, dest = cross_device_upload - dest.parent.mkdir(parents=True) - elsewhere = tmp_path / "elsewhere.png" - elsewhere.write_bytes(twin) - os.link(elsewhere, dest) - - _upload(temp) - - assert elsewhere.read_bytes() == twin - assert not os.path.samefile(elsewhere, dest) - assert dest.read_bytes() == _CONTENT - - -def test_failed_cross_device_copy_leaves_no_partial_file( - mock_create_session, cross_device_upload, monkeypatch -): - temp, dest = cross_device_upload - - def disk_full(src, dst): - Path(dst).write_bytes(_CONTENT[:5]) - raise OSError(errno.ENOSPC, "No space left on device") - - monkeypatch.setattr(ingest_module.shutil, "copyfile", disk_full) - - with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): - _upload(temp) - - assert not dest.exists() - assert not temp.exists() - - def test_other_move_errors_raise_without_copying( mock_create_session, cross_device_upload, monkeypatch ): From f089d822f82c3e1e1031068feb78e22c8f2606cf Mon Sep 17 00:00:00 2001 From: Simon Pinfold Date: Sun, 4 Oct 2026 15:16:13 -0700 Subject: [PATCH 14/15] test(assets): pin that a failed cross-volume copy fails the upload and registers nothing --- .../services/test_upload_cross_device.py | 21 +++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/tests-unit/assets_test/services/test_upload_cross_device.py b/tests-unit/assets_test/services/test_upload_cross_device.py index c869c22a8e2..d9e63cab69d 100644 --- a/tests-unit/assets_test/services/test_upload_cross_device.py +++ b/tests-unit/assets_test/services/test_upload_cross_device.py @@ -83,6 +83,27 @@ def test_cross_device_upload_replaces_other_bytes_at_the_destination( assert dest.read_bytes() == _CONTENT +def test_failed_cross_device_copy_fails_the_upload_and_registers_nothing( + mock_create_session, cross_device_upload, monkeypatch +): + temp, dest = cross_device_upload + + def disk_full(src, dst): + Path(dst).write_bytes(_CONTENT[:5]) + raise OSError(errno.ENOSPC, "No space left on device") + + monkeypatch.setattr(ingest_module.shutil, "copyfile", disk_full) + + with pytest.raises(RuntimeError, match="failed to move uploaded file into place"): + _upload(temp) + + assert not temp.exists() + with mock_create_session() as session: + assert session.scalars( + select(AssetContent).where(AssetContent.path == str(dest)) + ).first() is None + + def test_other_move_errors_raise_without_copying( mock_create_session, cross_device_upload, monkeypatch ): From 5f77928fab6e0e88de8e31b9d2656fc88fcbf238 Mon Sep 17 00:00:00 2001 From: Simon Pinfold Date: Sun, 4 Oct 2026 15:24:15 -0700 Subject: [PATCH 15/15] refactor(assets): drop the unused verified stat binding in uploads --- app/assets/services/ingest.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/app/assets/services/ingest.py b/app/assets/services/ingest.py index bb2220634ec..bbc5665b30e 100644 --- a/app/assets/services/ingest.py +++ b/app/assets/services/ingest.py @@ -456,7 +456,7 @@ def upload_from_temp_path( user_metadata = user_metadata or {} try: - digest, verified_stat = _snapshot_hash_with_retry(temp_path) + digest, _ = _snapshot_hash_with_retry(temp_path) except UploadUnstableError: _remove_temp_path(temp_path) raise