Skip to content

fix(episode): stage frame cache in unique tempdir and publish atomically - #635

Open
sivasurya05 wants to merge 1 commit into
Hebbian-Robotics:mainfrom
sivasurya05:fix/frame-cache-locking
Open

sivasurya05 wants to merge 1 commit into
Hebbian-Robotics:mainfrom
sivasurya05:fix/frame-cache-locking

Conversation

@sivasurya05

@sivasurya05 sivasurya05 commented Sep 26, 2026 •

Copy link
Copy Markdown

Motivation & Context

Under concurrent workloads sharing an explicit workdir (e.g., multi-worker training pipelines, PyTorch DataLoader(num_workers > 1), or prefetch), Episode.frames() and Episode.frames_at_indices() suffered from race conditions.

Problems

  1. Deterministic .tmp directory collisions: Both methods staged extracted frames using fixed directory names (frames_{label}.tmp and <output_dir>.tmp). Concurrent workers extracting frames for the same cache key could collide, delete, or overwrite each other's in-flight files.

  2. Completed cache deletion race: In frames_at_indices(), checking for missing frames and calling shutil.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:

tempfile.mkdtemp(
    dir=self.workdir,
    prefix=f"{output_dir.name}.",
    suffix=".tmp",
)

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:

try:
    staging_dir.rename(output_dir)
except OSError:
    # Another worker may have won the race.
    if not output_dir.exists():
        raise
finally:
    if staging_dir.exists():
        shutil.rmtree(staging_dir, ignore_errors=True)
  • The first caller to successfully publish establishes the completed cache.
  • If another caller has already published the cache, the losing caller discards its own staging directory and uses the existing cache.
  • Genuine filesystem errors, such as permission issues or a missing directory, are re-raised when the destination does not exist.
  • Temporary staging directories are cleaned up after publication or failure.

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:

RuntimeError(f"Incomplete frame cache at {output_directory}")

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: Uses threading.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 raises RuntimeError instead 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

uv run ruff check --fix
uv run ruff format
uv run ty check

Targeted Regression Tests

uv run pytest tests/test_episode.py -q -n 4

Result: 6 passed.

Full Test Suite

uv run pytest -q -n 4

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

@github-actions

Copy link
Copy Markdown

👋 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! 🙌

@greptile-apps

greptile-apps Bot commented Sep 26, 2026 •

Copy link
Copy Markdown
Contributor

RetriggerConfidence Score: 4/5

[Medium risk] Changes how frame extraction staging and caching work.

The PR does not yet appear safe to merge because an incomplete indexed cache still cannot be repaired.

Findings

  1. P1 Incomplete indexed caches go unrepaired ▶
  2. P2 Concurrent misses duplicate extraction ▶
Summary

This PR stages frame extraction in unique temporary directories and publishes completed caches by rename. It also adds concurrent-extraction and failure-cleanup tests. The latest changes defer incomplete-cache checks until after publication, avoiding a race with another worker.

Reviews (8) · Last reviewed commit: "fix(episode): stage frame caches uniquel..."

Comment thread src/hflow/episode.py Outdated
Comment thread src/hflow/episode.py Outdated
Comment thread src/hflow/episode.py Outdated
Comment thread src/hflow/episode.py Outdated
@sivasurya05
sivasurya05 force-pushed the fix/frame-cache-locking branch from aede263 to 120bbf9 Compare September 26, 2026 17:35

@kstonekuan kstonekuan left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 .tmp name.
  • Publish with a rename. If the output directory already exists, another caller finished first: remove your staging directory and use theirs.
  • Drop the rmtree of an existing output_directory in frames_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.

Comment thread src/hflow/episode.py
capture_output=True,
text=True,
check=False,
if not output_directory.exists():

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 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.

Comment thread src/hflow/episode.py
@sivasurya05 sivasurya05 changed the title fix(episode): coordinate frame cache extraction with filelock and safe staging fix(episode): stage frame cache in unique tempdir and publish atomically Sep 28, 2026
@sivasurya05
sivasurya05 force-pushed the fix/frame-cache-locking branch from ccee591 to 584cdfe Compare September 28, 2026 11:45
Comment thread src/hflow/episode.py Outdated
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.
@sivasurya05
sivasurya05 force-pushed the fix/frame-cache-locking branch from 584cdfe to 23a596b Compare September 28, 2026 12:14

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Bug]: Episode.frames and frames_at_indices race on deterministic .tmp staging paths under concurrency

2 participants