fix: finish cancelled uploads safely and isolate regression fixtures

This commit is contained in:
Palash Debnath
2026-09-28 16:21:17 -07:00
parent 60e33ea479
commit b1814197d7
5 changed files with 90 additions and 13 deletions
+24 -4
View File
@@ -818,9 +818,7 @@ async def dub_upload(
pass # Best effort: keep the original upload error.
raise
try:
await asyncio.to_thread(_stream_upload_to_disk)
except OSError as exc:
def _discard_upload():
try:
os.unlink(video_path)
except OSError:
@@ -828,7 +826,29 @@ async def dub_upload(
try:
os.rmdir(job_dir)
except OSError:
pass # Best effort: keep the original upload error.
pass # Best effort: only remove our empty reserved directory.
copying = asyncio.create_task(asyncio.to_thread(_stream_upload_to_disk))
try:
await asyncio.shield(copying)
except asyncio.CancelledError:
# A cancelled await cannot stop a file-copy thread. Keep its input open
# until it stops, then remove only this reserved upload's partial copy.
while not copying.done():
try:
await asyncio.shield(copying)
except asyncio.CancelledError:
continue
except Exception:
break
try:
copying.result()
except Exception:
pass # Preserve the caller's cancellation if the copy also failed.
_discard_upload()
raise
except OSError as exc:
_discard_upload()
if exc.errno == errno.ENOSPC or getattr(exc, "winerror", None) == 112:
raise _dub_upload_disk_error() from exc
raise
+2 -2
View File
@@ -1420,7 +1420,7 @@ async def ingest_pipeline(
else:
os.replace(extract_hq_path, audio_hq_path)
except Exception as e_hq: # noqa: BLE001 — quality upgrade, never fatal
logger.warning("HQ audio extraction errored (%s) — falling back", log_safe(e_hq))
logger.warning("HQ audio extraction errored (%s) — falling back", type(e_hq).__name__)
_discard_partial_audio(extract_hq_path)
audio_hq_path = None
except asyncio.CancelledError:
@@ -1428,7 +1428,7 @@ async def ingest_pipeline(
raise
except Exception as e:
_discard_partial_audio(extract_path, extract_hq_path)
logger.error("Extract failed for job %s: %s", log_safe(job_id), log_safe(e))
logger.error("Extract failed for job %s: %s", log_safe(job_id), type(e).__name__)
if isinstance(e, failure.InvalidMediaFileError):
await asyncio.to_thread(_discard_invalid_source_copy, job_dir, video_path)
yield prep_event("error", **failure.build_failure(e, stage="extract"))
+2 -4
View File
@@ -728,8 +728,7 @@ def raise_for_audio_extract_failure(stderr, path: str) -> None:
if already or is_no_audio_stream_stderr(stderr) or has_audio_stream(path) is False:
text = stderr.decode("utf-8", errors="replace") if isinstance(stderr, bytes) else str(stderr or "")
logger.info(
"Audio decode of %s failed because it has no audio stream: %s",
log_safe(os.path.basename(str(path))), log_safe(text[-500:]),
"Audio decode failed because the source has no audio stream"
)
raise NoAudioTrackError()
text = stderr.decode("utf-8", errors="replace") if isinstance(stderr, bytes) else str(stderr or "")
@@ -740,8 +739,7 @@ def raise_for_audio_extract_failure(stderr, path: str) -> None:
"invalid data found when processing input",
)):
logger.info(
"Audio decode of %s failed because its container is unreadable: %s",
log_safe(os.path.basename(str(path))), log_safe(text[-500:]),
"Audio decode failed because the source container is unreadable"
)
raise InvalidMediaFileError()
+3
View File
@@ -337,3 +337,6 @@ directory remains unchanged.
Extraction publishes each WAV only after FFmpeg succeeds, preserving completed
audio if a later import fails validation or runs out of space.
Cancelling an upload waits for its copy worker to stop before closing the input
and clearing the reserved job, so the same upload can be retried safely.
+59 -3
View File
@@ -10,12 +10,11 @@ from pathlib import Path
import pytest
from core import failure
from services import dub_pipeline
from services.ffmpeg_utils import raise_for_audio_extract_failure, validate_media_source
def test_zeroed_header_is_rejected_without_invoking_ffmpeg(tmp_path):
from core import failure
from services.ffmpeg_utils import validate_media_source
source = tmp_path / "original.mkv"
source.write_bytes(b"\0" * 4096 + b"remaining data")
@@ -24,12 +23,15 @@ def test_zeroed_header_is_rejected_without_invoking_ffmpeg(tmp_path):
def test_valid_media_header_is_left_for_ffmpeg_to_probe(tmp_path):
from services.ffmpeg_utils import validate_media_source
source = tmp_path / "original.mkv"
source.write_bytes(bytes.fromhex("1a45dfa3") + b"matroska payload")
validate_media_source(str(source))
def test_ffmpeg_unreadable_container_gets_the_same_guidance(tmp_path, monkeypatch):
from core import failure
from services.ffmpeg_utils import raise_for_audio_extract_failure
source = tmp_path / "original.mkv"
source.write_bytes(b"nonzero but damaged header")
monkeypatch.setattr("services.ffmpeg_utils.has_audio_stream", lambda _: None)
@@ -42,6 +44,7 @@ def test_ffmpeg_unreadable_container_gets_the_same_guidance(tmp_path, monkeypatc
def test_zeroed_source_emits_actionable_extract_failure(tmp_path, monkeypatch):
from services import dub_pipeline
job_dir = tmp_path / "job"
job_dir.mkdir()
source = job_dir / "original.mkv"
@@ -64,6 +67,7 @@ def test_zeroed_source_emits_actionable_extract_failure(tmp_path, monkeypatch):
def test_failed_copy_cleanup_never_removes_external_source(tmp_path, monkeypatch):
from services import dub_pipeline
job_dir = tmp_path / "dub_jobs" / "job"
job_dir.mkdir(parents=True)
original = tmp_path / "original.mkv"
@@ -207,3 +211,55 @@ def test_failed_reingest_preserves_completed_audio(tmp_path, monkeypatch, fail_b
asyncio.run(collect())
assert {p.name: p.read_bytes() for p in job_dir.iterdir()} == {
"audio.wav": b"completed audio", "audio_hq.wav": b"completed audio"}
def test_cancelled_upload_waits_for_writer_before_cleanup(tmp_path, monkeypatch):
import threading
from fastapi import UploadFile
from api.routers import dub_core
job_dir = tmp_path / "job"
started, release = threading.Event(), threading.Event()
upload = UploadFile(file=io.BytesIO(b"test media"), filename="clip.wav", size=10)
monkeypatch.setattr(dub_core, "_safe_job_dir", lambda _: str(job_dir))
monkeypatch.setattr(dub_core.shutil, "disk_usage", lambda _: SimpleNamespace(free=1024**2))
def copy(source, target, length):
target.write(source.read(2))
started.set()
assert release.wait(5)
target.write(source.read())
monkeypatch.setattr(dub_core.shutil, "copyfileobj", copy)
async def run():
waiting = asyncio.Event()
shield = asyncio.shield
calls = 0
def observed(future):
nonlocal calls
calls += 1
if calls == 2: waiting.set()
return shield(future)
monkeypatch.setattr(asyncio, "shield", observed)
task = asyncio.create_task(dub_core.dub_upload(video=upload, job_id="job", input_type="audio", source_lang=None))
try:
assert await asyncio.to_thread(started.wait, 3)
task.cancel()
await asyncio.wait_for(waiting.wait(), 3)
assert not upload.file.closed
task.cancel()
finally:
release.set()
with pytest.raises(asyncio.CancelledError): await task
asyncio.run(run())
assert upload.file.closed
assert not job_dir.exists()
def test_invalid_media_log_omits_source_path(tmp_path, monkeypatch, caplog):
from core import failure
from services import ffmpeg_utils
source = tmp_path / "private-recording.mkv"
monkeypatch.setattr(ffmpeg_utils, "has_audio_stream", lambda _: None)
with caplog.at_level("INFO"), pytest.raises(failure.InvalidMediaFileError):
ffmpeg_utils.raise_for_audio_extract_failure(
f"EBML header parsing failed: {source}", str(source))
assert str(source) not in caplog.text
assert "private-recording" not in caplog.text