Conversation
Cache-aware feature-bufferer and encoder-cache slots were returned only on is_last. Callers that end streams with delete_state (notably the simulstream adapter) leaked slots until num_slots was exhausted. Add idempotent free_stream helpers and restore all slots on session reset, keeping the existing is_last free path unchanged. Fixes NVIDIA-NeMo#16309 Signed-off-by: Dundy Pasupuleti <dundysm@gmail.com>
|
pinging @naymaraq @nithinraok @artbataev @lilithgrigoryan for cache-aware streaming / asr review. happy to rebase if #16231 lands first (same files, different concern). |
|
Thanks for picking this up, @dundysm. I ran this branch (ff02c54) and
Your test file: 13 passed here; copied onto The last two rows are extra motivation: on #16308 adds def delete_state(self, stream_id: int) -> None:
state = self.get_state(stream_id)
if (
state is not None
and self.decoding_computer is not None
and self.decoding_computer.per_stream_biasing_enabled
):
release_auto_managed_stream_biasing(state, self.decoding_computer.biasing_multi_model)
self._free_cache_aware_slots(stream_id)
super().delete_state(stream_id)The buffered pipelines' per-stream bufferers (RNN-T and CTC, left out of scope here, and SALM's A question: |
kzos
left a comment
There was a problem hiding this comment.
Two inline notes on the lines they refer to; the measurements are in my comment above.
|
|
||
| def reset_session(self) -> None: | ||
| """Reset the state pool and restore all cache-aware feature-bufferer and encoder-cache slots.""" | ||
| self.context_manager.reset() |
There was a problem hiding this comment.
CacheAwareContextManager.reset() builds a new cache through get_initial_cache_state(num_slots) while the old one is still referenced, and pipeline.run() reaches this twice (through open_session() and close_session()). On CPU with nemotron-speech-streaming-en-0.6b and 256 slots I measured +1.6 GiB peak RSS and roughly 0.2–0.3 s per call here, where main allocates nothing (details in my comment above; not measured on GPU). The initial cache is all zeros for this model, and zeroing the existing tensors in place took roughly 0.1 s with no extra peak. Would in-place zeroing work for you here? The CTC pipeline's reset_session() makes the same call (+2.6 GiB and roughly 0.3–0.5 s with 1024 slots).
| self.bufferer.reset() | ||
| super().reset_session() | ||
|
|
||
| def delete_state(self, stream_id: int) -> None: |
There was a problem hiding this comment.
#16308 (mine) adds CacheAwareRNNTPipeline.delete_state() at this same spot, to release the stream's per-stream biasing model, so the two conflict whichever lands second. The combined method is in my comment above: release the biasing model while the state still exists, then _free_cache_aware_slots(), then super().delete_state(). If this lands first, I'll rebase #16308 onto it.
Zero encoder cache tensors in place on CacheAwareContextManager.reset when they already exist, matching bufferer reset and avoiding realloc peak RSS on open_session/close_session. Fold NVIDIA-NeMo#16308 release order into CacheAwareRNNTPipeline.delete_state (biasing while state exists, then slots, then super). Add unit test that reset zeros without reallocating. Signed-off-by: Dundy Pasupuleti <dundysm@gmail.com>
|
thanks @kzos for the thorough cpu verification and for confirming the fix on both rnnt and ctc, including the stream-id reuse / simulstream correctness rows. that extra motivation is helpful. agree on both follow-ups:
i pushed a follow-up commit for (1) and (2), kept the existing 13 cpu unit tests green, and added a small check that reset() zeros in place rather than reallocating. your offered cpu reuse test that catches a free_stream skipping zeroing is welcome if you want to share it; free_stream here already goes through _reset_slots so it should already zero. buffered per-stream bufferers / salm left to #16322 as you noted. |
What does this PR do ?
Frees cache-aware feature-bufferer and encoder-cache slots when a stream ends via
delete_state()or when a session is reset/closed/opened, so callers that never sendis_last=True(notablyNeMoStreamingPipelineAdapter) no longer exhaustnum_slots.Collection: ASR
cc @naymaraq @nithinraok @artbataev @lilithgrigoryan
Changelog
BatchedCacheFeatureBufferer.free_stream(stream_id): idempotent return of the stream's slot (no-op if already freed viais_lastinupdate()).BatchedCacheFeatureBufferer.reset(): clear maps, refillavailable_slots, zero feature buffer and audio bufferers.CacheAwareContextManager.free_stream(stream_id): membership-guarded free so double free cannot corruptfree_slots(unlike raw_reset_slots).CacheAwareRNNTPipelineandCacheAwareCTCPipeline: overridedelete_stateto free both slot pools, andreset_sessionto restore all slots (coversopen_session/close_sessionvia the base path).num_slots=2, no checkpoint / HF download.The problem
On
main, both slot pools free a stream only when a request carriesis_last=True:BatchedCacheFeatureBufferer.updateappends toslots_to_freeonly forframe.is_lastCacheAwareContextManager.reset_slotsfrees only streams whose eos flag is true (eos_flagscome fromrequest.is_last)BasePipeline.delete_stateand session open/close only clear_state_pool. The simulstream adapter always sendsis_last=Falseand ends withdelete_state(), so afternum_slotsstreams the next allocate raisesRuntimeError: No free slots available.close_session()did not restore the pools either.Approach (design-preserving)
Keep the intentional
is_lastfree path for graceful EOU / decoding / bias release. Extend abrupt teardown:delete_state(stream_id)frees both slots if still mapped (idempotent after a prioris_lastfree).reset_session(and thereforeopen_session/close_session) rebuilds both pools.Do not flip the adapter to
is_last=Trueas the sole fix: that would change EOU /keep_all_outputs/ hypothesis-reset semantics and still leave aborted streams / session boundaries leaking.Why not only adapter
is_last=TrueReporter already flagged that option and has not measured output impact. Slot free via
delete_statematches how the adapter documents end-of-stream today. Session reset remains necessary either way.Tests
tests/collections/asr/inference/test_cache_aware_slot_release.py(CPU, no HF):free_streamreturns capacity and is idempotentresetrestores a fully leaked poolis_lastpath viareset_slots(..., eos=True)still works; followingfree_streamdoes not over-fill the queuedelete_staterecover free counts (issue repro shape withnum_slots=2), for RNNT and CTCdelete_stateafter a simulatedis_lastfree is a no-op on the queuesclose_session/open_sessionrestore leaked slotsBasePipeline.delete_statealone still leaks (documents the bug)Local run:
Related, deliberately not in this PR
delete_state): sibling leak; leave to that PR.is_last-only free; expect rebase if it lands first. This PR adds helpers beside the hotupdatepath and does not rewrite it.BatchedAudioBuffererper-stream dict growth: soft leak, out of scope.Before your PR is "Ready for review"
Pre checks:
PR Type:
Additional Information