fix(episode): stage frame cache in unique tempdir and publish atomically - #635
sivasurya05 wants to merge 1 commit into
Conversation
|
👋 Hi @sivasurya05 — thank you so much for your first contribution to HFlow! A maintainer will review your pull request as soon as possible. In the meantime:
💡 Tip: one open pull request per contributor at a time. Before starting an issue, check its sidebar for an assignee or a linked pull request. Either one means somebody is already on it; everything else is fair game. We are excited to have you here and appreciate your help making the project better! 🙌 |
|
aede263 to
120bbf9
Compare
kstonekuan
left a comment
There was a problem hiding this comment.
Thanks, the race is real when several episodes share one explicit workdir. The default workdir is private to each Episode, so that is the case to fix.
It doesn't need a lock. video() already solves the same problem: write_access_units_to_mp4 publishes atomically, so a finished file at the final path is always complete. Do the same for both frame caches:
- Stage in a unique directory per call (
tempfile.mkdtemp(dir=self.workdir, prefix=...)) instead of the fixed.tmpname. - Publish with a rename. If the output directory already exists, another caller finished first: remove your staging directory and use theirs.
- Drop the
rmtreeof an existingoutput_directoryinframes_at_indices. With atomic publishing it is either complete or absent.
That is about 20 lines in episode.py. For the test, two threads extracting the same frames into one workdir, both getting complete results and leaving no staging directories behind, is enough.
120bbf9 to
ccee591
Compare
| capture_output=True, | ||
| text=True, | ||
| check=False, | ||
| if not output_directory.exists(): |
There was a problem hiding this comment.
Incomplete indexed caches go unrepaired If an indexed cache directory exists but one of its frame files has been removed, this check skips extraction and returns a path to the missing file. A caller reading that frame then gets
FileNotFoundError; the previous per-file check rebuilt incomplete caches.
ccee591 to
584cdfe
Compare
Stage frame extraction in unique temporary directories per caller via tempfile.mkdtemp rather than a fixed .tmp name. Publish completed frames with an atomic rename, re-raising OSError unless the output directory was already published by a winning caller. Disallow returning broken paths by rejecting incomplete cache directories with RuntimeError. Add regression tests for barrier-enforced concurrency, incomplete cache rejection, rename error propagation, and failure cleanup.
584cdfe to
23a596b
Compare
Motivation & Context
Under concurrent workloads sharing an explicit
workdir(e.g., multi-worker training pipelines, PyTorchDataLoader(num_workers > 1), orprefetch),Episode.frames()andEpisode.frames_at_indices()suffered from race conditions.Problems
Deterministic
.tmpdirectory collisions: Both methods staged extracted frames using fixed directory names (frames_{label}.tmpand<output_dir>.tmp). Concurrent workers extracting frames for the same cache key could collide, delete, or overwrite each other's in-flight files.Completed cache deletion race: In
frames_at_indices(), checking for missing frames and callingshutil.rmtree(output_directory)could allow a lagging worker to delete another worker's freshly completed cache.These issues could result in missing frames, corrupted cache directories, and failed concurrent extraction workflows.
Solution
This change adopts the atomic publication pattern already established in
video()(write_access_units_to_mp4) to make frame cache creation safe under concurrent workloads.1. Unique Staging Directories
Each caller extracts frames into its own unique temporary directory using:
This ensures that concurrent workers never interfere with each other's staging directories.
2. Atomic Cache Publication
Once extraction completes, the staging directory is published to the final cache path using an atomic rename:
3. Strict Incomplete Cache Detection
In
frames_at_indices(), the cache is validated to ensure all expected frame files exist.If an existing cache directory is missing any expected frames, the method raises:
This prevents returning invalid frame paths or performing unsafe, uncoordinated cache deletions.
Regression Tests Added
The following regression tests were added to
tests/test_episode.py:test_concurrent_frame_extractions_share_workdir: Usesthreading.Barrier(2)around subprocess execution to ensure both workers encounter the cache miss concurrently. Verifies that both receive identical completed frames, all returned paths exist, and no temporary directories remain.test_frames_at_indices_rejects_incomplete_cache: Verifies that a cache directory missing expected frame files raisesRuntimeErrorinstead of returning nonexistent paths.test_rename_failure_propagates_and_cleans_up: Ensures genuine filesystem errors during rename are propagated and temporary staging directories are cleaned up.test_failed_ffmpeg_extraction_cleans_up_and_allows_retry: Verifies that a failed FFmpeg extraction leaves no cache or staging directories behind and that a subsequent retry succeeds.Validation
All checks passed successfully.
Quality & Style Checks
Targeted Regression Tests
Result: 6 passed.
Full Test Suite
Result: 73 passed.
Concurrency Stress Verification
Repeated execution of concurrency and failure tests 10 times.
Result: 100% pass rate.
Summary
This change eliminates deterministic staging directory collisions, prevents concurrent workers from deleting completed frame caches, and ensures incomplete caches are detected explicitly.
Frame extraction now uses isolated staging directories and atomic publication, improving reliability for concurrent workloads sharing a common
workdir.Closes #634