Skip to content

fix(asr): drop a buffered stream's bufferer when its state is deleted - #16322

Open
kzos wants to merge 1 commit into
NVIDIA-NeMo:mainfrom
kzos:fix/buffered-streaming-release-bufferer-on-delete-state
Open

kzos wants to merge 1 commit into
NVIDIA-NeMo:mainfrom
kzos:fix/buffered-streaming-release-bufferer-on-delete-state

Conversation

@kzos

@kzos kzos commented Sep 30, 2026

Copy link
Copy Markdown
Contributor

What does this PR do ?

Makes delete_state() in the buffered streaming pipelines (RNN-T, CTC and SALM) drop the stream's per-stream bufferer, and open_session() / close_session() drop all of them. A stream that ends without an is_last=True request no longer leaves its bufferer behind, and a later stream on the same id starts from an empty buffer.

Collection: ASR

Part of #16309. #16321 frees the cache-aware pipelines' slots and leaves the buffered pipelines' per-stream bufferers out of scope; this PR covers those. It does not touch the cache-aware pipelines, context_manager.py or cache_feature_bufferer.py. #16321's description links the issue with a closing keyword, so the issue will be closed when #16321 merges; this PR is the buffered part.

Changelog

  • BasePipeline.delete_state(): when the pipeline's bufferer is a BatchedAudioBufferer or BatchedFeatureBufferer (what init_bufferer_for_buffered_streaming() builds for BufferedRNNTPipeline and BufferedCTCPipeline), also call its existing rm_bufferer(stream_id), the call the is_last path makes.
  • BasePipeline.reset_session() (called by open_session() and close_session()): under the same condition, call the bufferer's existing reset(). Its docstring already reads "Reset the frame buffer and internal state pool".
  • BufferedSALMPipeline: new delete_state() and reset_session() overrides doing the same for its audio_bufferer (a BatchedIncrementalAudioBufferer), which it does not store as bufferer.
  • tests/collections/asr/inference/test_buffered_bufferer_release.py: 5 tests, 23 cases.

The problem

On main (00278b0) each of the three buffered bufferers keeps a dict with one bufferer per stream id. It adds the entry on the stream's first request and removes it only while processing a request with is_last (audio_bufferer.py L129-130, feature_bufferer.py L158-159, incremental_audio_bufferer.py L171-172). On main, BasePipeline.delete_state() only drops the state (base_pipeline.py L141-144), reset_session() only clears the state pool (L158-160), and none of the buffered pipelines overrides either.

The in-tree simulstream adapter ends every stream without is_last: process_chunk() sends is_last=False (simulstream_pipeline_adapter.py L311) and end_of_stream() calls delete_state() (L417). It forces request_type='frame' (L137), so on these pipelines it uses BatchedAudioBufferer.

This has two effects:

  1. Growth. The bufferer of each stream ended this way stays for the life of the pipeline. Closing and reopening the session does not remove it.
  2. Stream-id reuse, frame requests only. None of the three bufferers reads is_first, so a new stream on an id whose bufferer is still held continues that buffer: its first requests see the old stream's last audio as left context and a smaller left padding. A feature-buffer request carries its whole buffer, which FeatureBufferer.update() copies in (feature_bufferer.py L84), so there the leftover bufferer only costs memory.

The adapter does not reuse ids within one instance (clear() increments stream_id). But every instance starts at stream_id = 0 (L96), and all instances share one class-level pipeline (load_model() returns early once it is built, L118-119). So in a process that creates a second adapter instance, its first stream starts on the id of the first instance's first stream. I do not know whether simulstream creates more than one instance per process; I measured what happens when it does.

Measurements

All runs were on CPU in float32 with local checkpoints: stt_en_conformer_transducer_small in buffered_rnnt.yaml and stt_en_conformer_ctc_small in buffered_ctc.yaml (both 4.8 s chunks, 8.0 s buffers). I overrode only model, device, dtype, ITN/NMT (off) and log level, plus the settings stated below.

  • Growth. After 300 streams ended with delete_state() and no is_last request, main held 300 per-stream bufferers, and still 300 after close_session(). That held for RNN-T and CTC, through the adapter (frame requests) and through transcribe_step() directly (frame and feature-buffer requests). With this PR: 0.
  • Stream-id reuse. Over 12 pairs of LibriSpeech test-other utterances with frame requests, I reused a stream id after delete_state(). On main the second stream's text changed in 3 of 12 pairs (RNN-T) and 4 of 12 (CTC), and its segments in 8 and 12. With this PR nothing changed. With feature-buffer requests nothing changed on either side.
  • Streams that end with is_last. pipeline.run() outputs are identical between main and this PR in all 8 configurations described below (16 runs).
Bufferers held after 300 streams

One 4.8 s request of synthetic noise (0.01 * N(0, 1)) per stream, 300 streams. The adapter runs use a stub for the two names the adapter imports from simulstream, which is not installed here; the adapter and pipeline code that ran is NeMo's own.

driven by request type main: after 300 streams / after close_session() this PR
adapter (process_chunk(), end_of_stream(), clear()), RNN-T and CTC frame 300 / 300 0 / 0
transcribe_step() then delete_state(), RNN-T and CTC frame 300 / 300 0 / 0
transcribe_step() then delete_state(), RNN-T and CTC feature_buffer 300 / 300 0 / 0
transcribe_step(), streams left open (RNN-T frame, CTC feature_buffer) 300 / 300 300 / 0
control: every request with is_last=True, RNN-T and CTC both 0 / 0 0 / 0

In the transcribe_step() runs on main, the count also stayed at 300 after a following open_session(). The held buffer tensors alone were 512,000 bytes per stream with frame requests (153,600,000 for 300) and 256,320 bytes per stream with feature-buffer requests (76,896,000 for 300). I did not measure Python object overhead or process memory.

I also reran the adapter harness from #16308 on buffered RNN-T: 300 streams, each with a biasing request. These are the counts after the 300th stream, and again after close_session():

active biasing models per-stream bufferers
main 300 300
this PR 300 0
this PR + #16308 0 0
this PR + #16308 + #16321 0 0
Stream-id reuse

I took 12 pairs of consecutive LibriSpeech test-other utterances (X, Y). For each pair I compared two runs of Y:

  • Y streamed alone in a fresh session.
  • Y streamed on stream id 7 right after the first (up to 3) requests of X on id 7, sent with is_last=False and followed by delete_state(7).

The requests came from the pipeline's own request generator, as pipeline.run() sends them. Y alone gave the same output twice in all 12 pairs.

pipeline request type main: pairs where Y's text / segments (text, start, end, conf) differ this PR
RNN-T frame 3 / 8 of 12 0 / 0
CTC frame 4 / 12 of 12 0 / 0
RNN-T feature_buffer 0 / 0 0 / 0
CTC feature_buffer 0 / 0 0 / 0

On main, start times changed in 8 RNN-T pairs, and confidences changed in all 12 CTC pairs. The text changed in 7 runs:

  • In 6, Y's transcript gained a prefix ending in X's last words ("silent", "had a pretty good stock", "people").
  • In the seventh (CTC), one word inside Y changed: "he could do that" became "he could not do that". Here the changed output happens to match the reference text.

An RNN-T example, where X (7902-96591-0008) ends with "... had a pretty good stock":

Y alone : and all comes of dressing up in this stupid way like a rough fisher lad
after X : had a pretty good stock and all comes of dressing up in this stupid way like a rough fisher lad

Through the adapter, I streamed the first 12 utterances in two ways: with one instance (stream ids not used before), and with a new instance per utterance (stream id 0 each time). On main the transcripts of the two ways differed in 4 of 12 (RNN-T) and 5 of 12 (CTC). With this PR they differed in 0 of 12, and the new-instance transcripts equal main's one-instance transcripts.

Behaviour that should not change, and does not

I ran pipeline.run() on the first 12 LibriSpeech test-other utterances with batch_size=4, every stream ending with an is_last request. The following are identical between main and this PR for every stream:

  • the text
  • the segments or words, with their start and end times and confidences
  • the finalization delays

That covers RNN-T and CTC, frame and feature-buffer requests, and segment and word granularity: 8 configurations, 16 runs, 12 of 12 streams each. No bufferer was held after any of these runs on either side. On this path the new calls find nothing to remove: the is_last request has already removed the stream's bufferer when transcribe_step() calls delete_state() (base_pipeline.py L304-305 in this PR, L299-300 on main). On the merge of this PR with #16308 and #16321, the four segment-level configurations also match main.

Why BasePipeline, and composition with #16308, #16321 and #15827

BasePipeline builds these two bufferers in init_bufferer_for_buffered_streaming() (base_pipeline.py L442 and L450 in this PR, L437 and L445 on main), which the BufferedRNNTPipeline and BufferedCTCPipeline constructors call. Its reset_session() docstring already covers "the frame buffer". So I drop the bufferers in BasePipeline as well. Subclass overrides that end with super(), such as #16308's BufferedRNNTPipeline.delete_state() and #15827's BufferedRNNTPipeline.close_session(), reach the change without edits.

The isinstance check keeps BasePipeline off every other bufferer, in particular the cache-aware BatchedCacheFeatureBufferer, which #16321 handles in its own overrides. test_base_pipeline_leaves_other_bufferers_alone pins that.

BufferedSALMPipeline keeps its bufferer as audio_bufferer (buffered_salm_pipeline.py L105), so it gets its own two small overrides. If you would rather have one pattern, for example overrides in each buffered pipeline or a small hook they implement, I am happy to change it.

I merged this branch locally with each PR's current head (not pushed) and ran the PR test files and tests/collections/asr/inference/:

combination merge PR test files tests/collections/asr/inference/
main 167 passed, 7 skipped
this PR 23 passed 190 passed, 7 skipped
+ #16308 (ae69a65) clean 27 passed 194 passed, 7 skipped
+ #16321 (ff02c54) clean 36 passed 203 passed, 7 skipped
+ #15827 (0392df8, its diff applied) clean 23 passed 190 passed, 7 skipped

#16308 and #16321 both add CacheAwareRNNTPipeline.delete_state(), so they conflict with each other in cache_aware_rnnt_pipeline.py, with or without this branch, which does not touch that file. Whichever of the two lands second needs a rebase. I resolved that conflict by keeping both bodies. The three-way merge then gives the same test results as #16308 + #16321 without this branch, plus this PR's 23 passing cases.

The diffs of #16231, #16191 and #15940 (which change base_pipeline.py) apply cleanly in either order with this change, as do #15186's hunks in buffered_salm_pipeline.py. #15827, #16308 and #16321 do not touch the files changed here. #16158 targets another base branch. Its base_pipeline.py hunks do not apply to main either, and they start at L211, away from the lines changed here.

Tests

The tests are CPU-only and need no checkpoint. They build each buffered pipeline without a model. Its model step is replaced with the bufferer call that each real step begins with, recording the buffer and padding the model would get: buffered_rnnt_pipeline.py L416/L430, buffered_ctc_pipeline.py L310/L324, and buffered_salm_pipeline.py L184 in this PR (L174 on main). transcribe_step(), delete_state(), the sessions and the bufferers are the real ones. For RNN-T and CTC, the bufferers are the ones init_bufferer_for_buffered_streaming() builds.

  • test_delete_state_drops_bufferer_of_stream_ended_without_is_last (5 cases: RNN-T and CTC with frame and feature-buffer requests, and SALM): three streams dropped with delete_state() leave no bufferer, and an unknown id is ignored.
  • test_delete_state_drops_only_the_deleted_streams_bufferer (5 cases): delete_state() after an is_last request does nothing. The stream still open keeps its bufferer and gets the same buffer on its next request as on a pipeline where nothing was deleted.
  • test_close_and_open_session_drop_all_bufferers (5 cases).
  • test_reused_stream_id_starts_from_an_empty_buffer (6 cases: the three frame configurations, with the stream ended by delete_state() or by close_session() + open_session()): a new stream on the old id sees the same buffer and padding as on a fresh pipeline. Feature-buffer requests are not included, since they carry their whole buffer.
  • test_base_pipeline_leaves_other_bufferers_alone (2 cases, cache-aware RNN-T and CTC): BasePipeline.delete_state() and reset_session(), called directly, make no call on a bufferer of any other type. The test uses a Mock that records calls. It passes on main too, since main calls no bufferer there. Its job is to keep this change away from the cache-aware bufferer.
test_buffered_bufferer_release.py
without this change : 21 failed, 2 passed (the 2 are test_base_pipeline_leaves_other_bufferers_alone)
with this change    : 23 passed

tests/collections/asr/inference/
without this change : 167 passed, 7 skipped
with this change    : 190 passed, 7 skipped

The 7 skips need --with_downloads or nemo_text_processing, and I did not run them.

To check that the tests catch wrong fixes and not only a missing one, I swapped in eleven. Each failed some of the 23 cases. The two that loosen the isinstance check are caught only by test_base_pipeline_leaves_other_bufferers_alone.

Wrong fixes and the cases they fail
wrong fix cases failed
delete_state() unchanged 10
sessions unchanged 6
delete_state() dropping every stream's bufferer 4
handling only BatchedAudioBufferer 6 (the feature-buffer cases)
clearing the stream's buffer but keeping its entry 8
SALM not covered 5
SALM sessions unchanged 2
self.bufferer without getattr 5 (SALM)
del instead of rm_bufferer(), which raises for a stream that holds no bufferer 8
the isinstance check replaced by bufferer is not None 2 (only test_base_pipeline_leaves_other_bufferers_alone)
the isinstance check replaced by a hasattr check 2 (only test_base_pipeline_leaves_other_bufferers_alone)

black==24.10.0 and isort==5.13.2 report no changes on the changed files.

Limits

  • Serialized use only. The lifecycle calls have no locking, and I measured single-threaded use.
  • delete_state() is keyed by stream id only. A call that arrives after the id has been reused by a new stream deletes the new stream's state, as on main, and now also its bufferer.
  • SALM: I have no SALM checkpoint locally, so its change is covered by the unit tests only, not measured end to end.
  • The request generator that pipeline.run() uses for feature-buffer requests has its own BatchedAudioBufferer (multi_stream.py L291), which also removes a stream only at is_last. It lives for one run() call, and delete_state() does not reach it. It is not changed here.
  • The adapter is not changed.

Before your PR is "Ready for review"

Pre checks:

  • Make sure you read and followed Contributor guidelines
  • Did you write any new necessary tests?
  • Did you add or update any necessary documentation?
  • Does the PR affect components that are optional to install? (Ex: Numba, Pynini, Apex etc) — no

PR Type:

  • New Feature
  • Bugfix
  • Documentation

Additional Information

cc @naymaraq @lilithgrigoryan @artbataev @nithinraok

The buffered pipelines keep one audio or feature bufferer per stream id
and remove it only while processing a request with `is_last=True`.
`delete_state()` removed only the stream's state, and `close_session()`
and `open_session()` removed no bufferer, so a caller that ends its
streams with `delete_state()`, as the in-tree simulstream adapter does,
kept one bufferer per ended stream for the life of the pipeline. A new
stream on the id of such a stream also continued the old stream's audio
buffer.

`BasePipeline.delete_state()` now removes the stream's bufferer, and
`reset_session()` removes all of them, when the pipeline's bufferer is
the buffered pipelines' `BatchedAudioBufferer` or
`BatchedFeatureBufferer`, as a request with `is_last=True` does.
`BufferedSALMPipeline`, which keeps its bufferer as `audio_bufferer`,
does the same in its own `delete_state()` and `reset_session()`.
Removing the bufferer of a stream that holds none, such as one whose
`is_last` request already removed it or an unknown id, does nothing.
The cache-aware pipelines' bufferer is left alone.

Outputs of streams that end with `is_last` are unchanged: text,
segments or words with their times and confidences, and finalization
delays through `pipeline.run()` on 12 LibriSpeech test-other
utterances, for the buffered RNN-T and CTC pipelines with frame and
feature-buffer requests, on CPU.

Part of NVIDIA-NeMo#16309

Signed-off-by: Zaheer Sheriff K <zaheersheriff.k@gmail.com>
@copy-pr-bot

copy-pr-bot Bot commented Sep 30, 2026

Copy link
Copy Markdown

This pull request requires additional validation before any workflows can run on NVIDIA's runners.

Pull request vetters can view their responsibilities here.

Contributors can view more details about this message here.

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

1 participant