Skip to content

[Feature] Add a blocking replay buffer - #4419

Closed
vmoens wants to merge 7 commits into
mainfrom
replay-producer-admission
Closed

vmoens wants to merge 7 commits into
mainfrom
replay-producer-admission

Conversation

@vmoens

@vmoens vmoens commented Sep 17, 2026 •

Copy link
Copy Markdown
Contributor

Summary

  • add BlockingReplayBuffer, a bounded consuming ReplayBuffer subclass whose producers wait rather than overwrite sampleable records
  • use the existing consume_after_n_samples mechanism and storage capacity; each extend is admitted atomically as a complete batch
  • support timeout and cancellation on blocking writes so producer workers can exit cleanly
  • share ConsumingSampler masks and known length across processes, and rebuild its process-local free-list from shared state

Scope

This replaces producer-admission options on the ReplayBuffer base class. There are no drop or raise modes, configurable watermarks, new base statistics, wait_until_writable, Hydra fields, or remote, prioritized, TensorDict, and ensemble signature changes. Ordinary replay buffers retain their existing circular-overwrite behavior.

Remote services, prefetching, custom samplers, replay transforms, multidimensional storage, and compilable writers are intentionally outside this first specialized API.

Review fixes

  • the shared-process regression performs a successful worker write before the blocking write, then verifies that sampling in another process releases it
  • a separate shared-worker test verifies that cancel_event terminates a blocked write without modifying replay contents

Stack

Depends on #4408 for readiness notifications and consuming replay support. Replay-ratio limiting remains separate in #4418.

Validation

  • consuming-buffer, blocking, shared-growth, and cancellation tests: 24 passed
  • blocking consume/refill benchmark: passed
  • ufmt, flake8, and git diff --check: passed

@pytorch-bot

pytorch-bot Bot commented Sep 17, 2026 •

Copy link
Copy Markdown

🔗 Helpful Links

🧪 See artifacts and rendered test results at hud.pytorch.org/pr/pytorch/rl/4419

Note: Links to docs will display an error until the docs builds have been completed.

❗ 1 Active SEVs

There are 1 currently active SEVs. If your PR is affected, please view them below:

❌ 2 New Failures, 2 Unrelated Failures

As of commit 847bab8 with merge base 5757139 (image):

NEW FAILURES - The following jobs have failed:

BROKEN TRUNK - The following jobs failed but were present on the merge base:

👉 Rebase onto the `viable/strict` branch to avoid these failures

This comment was automatically generated by Dr. CI and updates every 15 minutes.

@meta-cla meta-cla Bot added the CLA Signed This label is managed by the Facebook bot. Authors need to sign the CLA before a PR can be reviewed. label Sep 17, 2026
@vmoens vmoens added the ci/optdeps Run the full tests-optdeps suite on this PR label Sep 17, 2026
@github-actions github-actions Bot added Feature New feature Documentation Improvements or additions to documentation Benchmarks rl/benchmark changes ReplayBuffers Trainers and removed Feature New feature ci/optdeps Run the full tests-optdeps suite on this PR labels Sep 17, 2026
@vmoens
vmoens force-pushed the replay-producer-admission branch from 29a1e46 to 59532cb Compare September 17, 2026 10:53
@github-actions github-actions Bot added the Feature New feature label Sep 17, 2026
@vmoens vmoens added the ci/optdeps Run the full tests-optdeps suite on this PR label Sep 17, 2026
@vmoens
vmoens force-pushed the replay-producer-admission branch from 59532cb to c1e22ab Compare September 17, 2026 10:53

@vmoens vmoens left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Review (COMMENT) — own PR, comment-only. (Delta on top of #4408 reviewed.)

The concurrency design is careful and I verified the hard cases:

  • Whole-write atomicity is real: the readiness condition is held across admission and the file write (_extend_unchecked runs inside the reservation context), so two producers cannot interleave a partial batch or jointly cross the watermark. The cost — all producer writes for an admission-controlled buffer serialize, and readiness waiters stall during a slow memmap write — is the honest price of the guarantee and the fast path for the default mode is untouched.
  • First-block notify_all isn't spurious: it's what lets an external observer watch producer_waiters rise, exactly as the shared-memory test does. Nice.
  • Lock order is uniform (condition → _replay_lock); _sample never takes the condition, pressure-set/clear happens under _replay_lock. No inversion found.
  • Ensemble rejection (can't atomically reserve across members), Ray block rejection (would wedge the owning actor), and distributed-transport rejection are the right default-deny postures, and they fail at construction with actionable messages.
  • Checkpoint parity across state_dict/dumps/loads, counter restore with waiters zeroed, and the 64-bit kinds are consistent. Backward-compat branch for the older two-counter context is present.

Findings, in descending order of importance:

  1. block on a non-consuming buffer is a latent deadlock that only the RST guide warns about. Physical occupancy never falls (sampling doesn't free capacity), so after the first reach of the high watermark a blocked producer only recovers via empty(). The parameter docstring currently says "waits for space" without that caveat — a one-line warning there ("without consumption, blocked producers resume only after empty()") would save a real user from this the first time; the guide's prose is one hop too far.

  2. wait_until_writable pollutes producer-behavior counters. It registers blocked_producer_calls and joins/leaves producer_waiters like a real producer. A monitoring loop polling writability will inflate the very metrics used to detect producer stalling. Either give advisory waits their own counters or document the shared accounting in wait_until_writable's docstring.

  3. The overwrite counter is an estimate for consuming buffers. occupancy is sampleable-based, but consumed-slot reuse in _pop_consumed_indices means actual live evictions can be lower than occupancy + N - capacity. Fine for telemetry, but the stats docstring should say "estimated".

  4. Cancelled drop_newest/raise writes are indistinguishable from policy drops. For non-block admissions a set cancel_event yields False → extend returns None, same as a drop; only block raises. Probably fine, but the extend docstring should say cancellation surfaces as None under the non-blocking policies.

  5. In TensorDictReplayBuffer.extend, the new if index is None: return None check comes after _set_index_in_td(tensordicts, index); harmless today only because _set_index_in_td early-returns on None. Move the guard above the call or add a comment; the current order will look broken to the next reader.

@vmoens

vmoens commented Sep 18, 2026

Copy link
Copy Markdown
Contributor Author

CI triage for c60cfac8b:

  • Branch-caused: the three CPU bulk jobs and optdeps fail in TestEnsemble::test_routed_write_and_conditional_update. The test compares the whole ReplayBufferEnsemble.stats() dict and still expects the pre-admission schema, while this PR correctly adds the producer-admission fields. The focused fix is to update that expected dict. I have not changed the branch because its worktree currently has an in-progress rebase conflict, which must be resolved before this test-only correction can be applied safely.
  • Unrelated: the Python 3.10/3.13 mp jobs have the repository-wide Trackio / CommitOperationAdd import mismatch; both GPU shard-2 jobs have the recurring CUDA Invalid stream capture status failure in TestProcessSlotTransport::test_server_batched_pass_on_cuda[True]; optdeps additionally has the unrelated stochastic SAC priority assertion in test_sac_prioritized_weights[2].

No remote branch mutation was made during this triage.

@vmoens

vmoens commented Sep 18, 2026

Copy link
Copy Markdown
Contributor Author

Follow-up: the branch-caused CI regression is fixed in d2d67b422.

ReplayBufferEnsemble.stats() correctly gained the producer-admission fields, but TestEnsemble::test_routed_write_and_conditional_update still compared against the old exact dictionary. The expected contract now includes overwrite/drop/blocking counters, waiter/pressure state, and admission watermarks.

I also completed the rebase onto the updated #4408 base, retaining both independent distributed blocking tests from the conflict.

Focused validation: 14 passed across the failing ensemble regression, the producer-admission suite, and the distributed blocking guards. A fresh CI run is now attached to the updated head.

Base automatically changed from replay-readiness-counters to main September 18, 2026 07:13
@vmoens
vmoens requested a review from theap06 September 18, 2026 07:24
@vmoens
vmoens force-pushed the replay-producer-admission branch from d2d67b4 to 4fa8c19 Compare September 18, 2026 07:42
@vmoens

vmoens commented Sep 18, 2026

Copy link
Copy Markdown
Contributor Author

CI follow-up for 4fa8c197d:

  • Branch-caused test race, fixed in 3cd26e03e: test_shared_replay_producer_unblocks_after_consumption started its five-second deadline at Process.start(), so under loaded bulk runners the deadline could expire before the child had bootstrapped. Both Python 3.10 and 3.13 failures crossed that deadline by less than 0.5 ms without observing a producer waiter. The test now synchronizes child readiness first, observes admission through the public stats() API, and always cleans up the child. The focused producer-admission suite passes (11 tests), and the multiprocess regression passes 20/20 repeated runs.
  • Unrelated Dreamer CUDA failure: test_discrete_actor[cuda-autocast] differs at one argmax element between repeated bfloat16 CUDA forwards. This PR changes only replay-buffer code/tests/docs and does not touch Dreamer or module code. The probability comparison immediately before the failed one-hot assertion passed; the one-element near-tie/argmax mismatch is numerical and unrelated to producer admission. A fresh run is attached to the new head.


def _producer_occupancy_locked(self) -> int:
if self._is_consuming():
return int(self._sampler._num_sampleable(self._storage))

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

per process, so a producer in another process never sees consumption and blocks forever. share the count or reject block admission on shared consuming buffers.

Comment thread test/rb/test_storages.py Outdated
assert rb.stats()["sample_calls"] == 1
assert rb.stats()["samples_returned"] == 2

def test_shared_replay_producer_unblocks_after_consumption(self):

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

only passes because the mask was allocated before spawn. add one successful worker write before the blocking one and it hangs.

while True:
if self._shutdown_requested():
raise RuntimeError(
"A shut down replay buffer cannot admit producer writes."

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

a blocked worker has no way out except being killed or this raising inside it. give the collector write a cancel event.

@vmoens vmoens changed the title [Feature] Add replay producer admission and backpressure [Feature] Add a blocking replay buffer Sep 18, 2026
@vmoens
vmoens requested a review from theap06 September 18, 2026 13:10
condition = self._readiness_condition
with condition:
while True:
if self._service_shutdown:

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

The learner side shutdown never reaches a worker here, and collector workers call extend with no cancel_event or timeout.

@vmoens
vmoens force-pushed the replay-producer-admission branch from 7026ca5 to 6c3ca32 Compare September 19, 2026 14:08
@vmoens
vmoens force-pushed the replay-producer-admission branch from 6c3ca32 to 847bab8 Compare September 19, 2026 14:23
@vmoens vmoens closed this Sep 22, 2026
@vmoens
vmoens deleted the replay-producer-admission branch September 22, 2026 14:24
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Benchmarks rl/benchmark changes ci/optdeps Run the full tests-optdeps suite on this PR CLA Signed This label is managed by the Facebook bot. Authors need to sign the CLA before a PR can be reviewed. Documentation Improvements or additions to documentation Feature New feature ReplayBuffers Trainers

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants