Conversation
🔗 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 SEVsThere are 1 currently active SEVs. If your PR is affected, please view them below: ❌ 2 New Failures, 2 Unrelated FailuresAs of commit 847bab8 with merge base 5757139 ( 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. |
29a1e46 to
59532cb
Compare
59532cb to
c1e22ab
Compare
vmoens
left a comment
There was a problem hiding this comment.
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_uncheckedruns 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_allisn't spurious: it's what lets an external observer watchproducer_waitersrise, exactly as the shared-memory test does. Nice. - Lock order is uniform (
condition → _replay_lock);_samplenever takes the condition, pressure-set/clear happens under_replay_lock. No inversion found. - Ensemble rejection (can't atomically reserve across members), Ray
blockrejection (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:
-
blockon 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 viaempty(). 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. -
wait_until_writablepollutes producer-behavior counters. It registersblocked_producer_callsand joins/leavesproducer_waiterslike 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 inwait_until_writable's docstring. -
The overwrite counter is an estimate for consuming buffers.
occupancyis sampleable-based, but consumed-slot reuse in_pop_consumed_indicesmeans actual live evictions can be lower thanoccupancy + N - capacity. Fine for telemetry, but the stats docstring should say "estimated". -
Cancelled
drop_newest/raisewrites are indistinguishable from policy drops. For non-block admissions a setcancel_eventyieldsFalse→extendreturnsNone, same as a drop; onlyblockraises. Probably fine, but theextenddocstring should say cancellation surfaces asNoneunder the non-blocking policies. -
In
TensorDictReplayBuffer.extend, the newif index is None: return Nonecheck comes after_set_index_in_td(tensordicts, index); harmless today only because_set_index_in_tdearly-returns onNone. Move the guard above the call or add a comment; the current order will look broken to the next reader.
|
CI triage for
No remote branch mutation was made during this triage. |
c60d3ff to
d2d67b4
Compare
|
Follow-up: the branch-caused CI regression is fixed in
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. |
d2d67b4 to
4fa8c19
Compare
|
CI follow-up for
|
|
|
||
| def _producer_occupancy_locked(self) -> int: | ||
| if self._is_consuming(): | ||
| return int(self._sampler._num_sampleable(self._storage)) |
There was a problem hiding this comment.
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.
| assert rb.stats()["sample_calls"] == 1 | ||
| assert rb.stats()["samples_returned"] == 2 | ||
|
|
||
| def test_shared_replay_producer_unblocks_after_consumption(self): |
There was a problem hiding this comment.
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." |
There was a problem hiding this comment.
a blocked worker has no way out except being killed or this raising inside it. give the collector write a cancel event.
| condition = self._readiness_condition | ||
| with condition: | ||
| while True: | ||
| if self._service_shutdown: |
There was a problem hiding this comment.
The learner side shutdown never reaches a worker here, and collector workers call extend with no cancel_event or timeout.
7026ca5 to
6c3ca32
Compare
6c3ca32 to
847bab8
Compare
Summary
BlockingReplayBuffer, a bounded consumingReplayBuffersubclass whose producers wait rather than overwrite sampleable recordsconsume_after_n_samplesmechanism and storage capacity; eachextendis admitted atomically as a complete batchConsumingSamplermasks and known length across processes, and rebuild its process-local free-list from shared stateScope
This replaces producer-admission options on the
ReplayBufferbase 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
cancel_eventterminates a blocked write without modifying replay contentsStack
Depends on #4408 for readiness notifications and consuming replay support. Replay-ratio limiting remains separate in #4418.
Validation
git diff --check: passed