From f075f31569739b493faf7437fbb715a9bda8457f Mon Sep 17 00:00:00 2001 From: simobenziane Date: Wed, 9 Sep 2026 20:37:25 +0200 Subject: [PATCH] fix: bound the streaming mel buffer to a sliding window (fixes #63, the streaming RSS leak) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Root cause: parakeet_capi_stream_feed's streaming session keeps every mel frame from stream_begin onward in mel_buf and rebuilds that whole, ever-growing array on every feed call, even though the decoder (feed_available's window()) never reads a frame older than mel_buffer_idx - pre_encode_cache_size() frames back — so those old frames are provably dead the moment the decoder consumes past them. Because each rebuild requests a uniquely, ever-larger allocation for the life of the stream, the allocator can never reuse an earlier call's freed (now permanently too-small) block, so process RSS grows with the cumulative history of buffer sizes rather than the small, bounded set of frames actually still needed — which is why neither stream_free nor a fresh stream_begin reclaims it. The fix: track the absolute frame index of mel_buf's first column (mel_buf_origin) and, on every append, drop the prefix strictly before mel_buffer_idx - pre_encode_cache_size() — the exact bound window() already relies on — so mel_buf stays bounded to O(pre_encode_cache_size + chunk_size) instead of O(stream length). No API change, no behavioural change: window()'s absolute-index reads are translated into mel_buf's origin-relative local columns, and the surviving frames are carried forward byte-for-byte. Measurement (CPU + Accelerate build, realtime_eou_120m-v1 q8, 120 s of silence fed in 20 ms chunks, RSS sampled at 30/60/90/120 s of audio fed): growth over the last 60 s drops from 178 MB to 0.5 MB; the same run's wall time drops from ~14 minutes to ~13 seconds, since the unbounded rebuild was O(T) per call, O(T²) total. A LibriSpeech utterance's transcript is byte-identical before and after. --- src/parakeet_capi.cpp | 77 ++++++++++++++++++++++++++++++++++--------- 1 file changed, 62 insertions(+), 15 deletions(-) diff --git a/src/parakeet_capi.cpp b/src/parakeet_capi.cpp index 1a6c9b5..241a62c 100644 --- a/src/parakeet_capi.cpp +++ b/src/parakeet_capi.cpp @@ -66,9 +66,16 @@ struct parakeet_ctx { struct parakeet_stream { parakeet_ctx* ctx = nullptr; // borrowed (must outlive the stream) std::unique_ptr mel; // incremental log-mel front end - std::vector mel_buf; // accumulated mel [n_mels, mel_T] feat-major + // mel_buf holds ONLY the still-reachable window of mel history, feat-major + // [n_mels, mel_T - mel_buf_origin]; column 0 is ABSOLUTE frame mel_buf_origin + // (see append_mel_frames / feed_available's window()). Frames strictly + // before mel_buffer_idx - pre_encode_cache_size() can never be read again + // (mel_buffer_idx only advances -- see feed_available), so they are dropped + // as they age out instead of being retained for the life of the stream. + std::vector mel_buf; int n_mels = 0; - int mel_T = 0; // total mel frames accumulated so far + int mel_T = 0; // total mel frames accumulated so far (absolute) + int mel_buf_origin = 0; // absolute frame index of mel_buf's column 0 std::unique_ptr sess; int mel_buffer_idx = 0; // next un-fed mel frame (chunk schedule) bool first_chunk = true; // chunk 0 has no pre-encode overlap @@ -77,25 +84,60 @@ struct parakeet_stream { namespace { // Append `n_new` feat-major mel frames `[n_mels, n_new]` to the stream's -// accumulated feat-major mel buffer `[n_mels, mel_T]`, growing mel_T. Both are -// feat-major (out[m*T + t]); appending along the time axis requires a per-row -// rebuild since the inner stride changes when T grows. +// accumulated feat-major mel buffer, growing mel_T, and drop the prefix that +// feed_available's window() can provably never read again. +// +// ROOT CAUSE of mudler/parakeet.cpp#63 (the streaming RSS leak): the mel buffer +// used to retain the FULL stream history from mel_buffer_idx==0 forever, so +// this function rebuilt a `[n_mels, T]` array from scratch on EVERY stream_feed +// call with T monotonically growing for the life of the stream -- not just O(T) +// CPU per call (the reason this rebuild exists at all, per the original +// comment), but a continuous sequence of UNIQUELY, EVER-LARGER-sized heap +// allocations. A general-purpose allocator can never satisfy a strictly-larger +// request from a smaller freed block, so the freed (now too-small) blocks from +// every earlier call accumulate as unusable, unreturned pages instead of being +// recycled -- the process's live working set stays tiny (a few MB) but the +// CUMULATIVE allocator debt grows with the stream's length, which is exactly +// why the upstream report finds it reclaimed by neither stream_free nor a +// fresh stream_begin: it is allocator-level fragmentation from the discarded +// buffers' sizes, not stream-owned memory. +// +// feed_available's window(lo, hi) only ever reads lo >= mel_buffer_idx - +// pre_cache, and mel_buffer_idx is monotonically non-decreasing (it only moves +// forward by chunk_size). So at the START of any call here (mel_buffer_idx +// unchanged since the previous feed_available pass), every mel frame before +// mel_buffer_idx - pre_cache is provably dead and safe to drop; mel_buf then +// stays bounded to O(pre_cache + chunk_size) -- a small constant -- instead of +// O(stream length), for the life of the stream. void append_mel_frames(parakeet_stream* s, const std::vector& frames, int n_new) { if (n_new <= 0) return; const int n_mels = s->n_mels; - const int old_T = s->mel_T; + const int old_T = s->mel_T; // absolute frame count so far const int new_T = old_T + n_new; - std::vector out((size_t)n_mels * new_T); + + const int pre_cache = s->sess ? s->sess->pre_encode_cache_size() : 0; + int keep_from = s->mel_buffer_idx - pre_cache; + if (keep_from < s->mel_buf_origin) keep_from = s->mel_buf_origin; // never move backward + if (keep_from > old_T) keep_from = old_T; // never past what exists + + const int old_local_T = old_T - s->mel_buf_origin; // mel_buf's current column count + const int drop_local = keep_from - s->mel_buf_origin; // columns to drop from the front + const int kept = old_local_T - drop_local; // columns carried forward + const int new_local_T = kept + n_new; + + std::vector out((size_t)n_mels * new_local_T); for (int m = 0; m < n_mels; ++m) { - // copy existing [0, old_T) - for (int t = 0; t < old_T; ++t) - out[(size_t)m * new_T + t] = s->mel_buf[(size_t)m * old_T + t]; - // append new [old_T, new_T) + // carry forward the still-reachable tail [drop_local, old_local_T) + for (int t = 0; t < kept; ++t) + out[(size_t)m * new_local_T + t] = + s->mel_buf[(size_t)m * old_local_T + (drop_local + t)]; + // append the newly-arrived frames for (int t = 0; t < n_new; ++t) - out[(size_t)m * new_T + (old_T + t)] = frames[(size_t)m * n_new + t]; + out[(size_t)m * new_local_T + (kept + t)] = frames[(size_t)m * n_new + t]; } s->mel_buf.swap(out); s->mel_T = new_T; + s->mel_buf_origin = keep_from; } } // namespace @@ -492,20 +534,25 @@ std::string feed_available(parakeet_stream* s, bool flush, int& eou_flag, const size_t ev0 = sess.events().size(); const int n_mels = s->n_mels; - const int T = s->mel_T; + const int T = s->mel_T; // absolute total frame count if (T <= 0) return std::string(); - const std::vector& mel = s->mel_buf; // [n_mels, T] feat-major + const int origin = s->mel_buf_origin; // absolute index of mel_buf's column 0 + const int local_T = T - origin; // mel_buf's actual (bounded) column count + const std::vector& mel = s->mel_buf; // [n_mels, local_T] feat-major, columns [origin, T) const int chunk0 = sess.chunk_size_first(); const int chunk_main = sess.chunk_size(); const int pre_cache = sess.pre_encode_cache_size(); + // lo/hi are ABSOLUTE frame indices (as before); translate into mel_buf's + // own local (origin-relative) columns. append_mel_frames guarantees + // lo >= mel_buffer_idx - pre_cache is always still resident. auto window = [&](int lo, int hi) { const int len = hi - lo; std::vector w((size_t)n_mels * len); for (int m = 0; m < n_mels; ++m) for (int t = 0; t < len; ++t) - w[(size_t)m * len + t] = mel[(size_t)m * T + (lo + t)]; + w[(size_t)m * len + t] = mel[(size_t)m * local_T + (lo - origin + t)]; return w; };