skip to content

A Dart StreamController republishing chat messages from a socket keeps growing in memory while its slow listener is paused — why, and how do you fix it?

level: seniorimportance: should knowfreq 30%

answer

  1. the controller buffers, it does not refuse
  2. isPaused ignored by the producer
  3. forward pause to the socket subscription
  4. addStream forwards pause and cancel
  5. stop the source in onCancel

basics

~20 s

A StreamController never refuses add: while its listener is paused, or before it listens, every event is queued. If the socket callback keeps calling add, the queue grows. Pause the socket subscription in onPause, or pipe it with addStream.

solid answer

~40 s

`StreamController.add` never blocks or rejects; while the single listener is paused (or before anyone listens) each event goes into the subscription's buffer, and `isPaused` reports `true`. A socket callback that calls `add` unconditionally therefore turns a slow consumer into an unbounded queue. The fix is to make the source feel the pause: subscribe to the socket in `onListen`, call `pause()` and `resume()` on that subscription from `onPause` and `onResume`, and cancel it in `onCancel`, so bytes stay in the OS socket buffer and TCP flow control slows the sender. When the controller just relays one stream, `controller.addStream(source)` or a `StreamTransformer` forwards pause, resume and cancel for you. Also check the opposite leak: after the listener cancels, `add` is silently discarded, so a producer that never looks at `onCancel` keeps working for nobody.

code

dart · 11 lines
dart
import 'dart:async';

Stream<List<int>> relay(Stream<List<int>> socket) {
  final controller = StreamController<List<int>>();
  controller.onListen = () {
    // Forwards data and errors; the controller pauses and resumes
    // the socket subscription whenever its own listener does.
    controller.addStream(socket).whenComplete(controller.close);
  };
  return controller.stream;
}

go deeper

for a junior

Know that a paused listener does not stop a StreamController from accepting events; they queue up.

for a middle

Explain the three states where add does not deliver immediately — before listen, while paused, after cancel — and what isPaused and hasListener report in each.

for a senior

Diagnose the leak from a growing retained set, and fix it by forwarding pause and cancel to the source through the hooks, addStream, or by returning a transformed stream instead of a controller.

for a principal

Decide where a pipeline may block, drop or buffer, and make every stream-producing component document whether it honours pause, so slow consumers cannot silently exhaust memory.

## The symptom A chat client parses frames from a socket and republishes them through a `StreamController<ChatMessage>`. The listener writes each message to a local database and occasionally pauses its subscription while a batch commits. Memory climbs steadily during long sessions, and after a pause the listener receives hundreds of messages at once. ## Why it happens A **`StreamController` has no way to say no.** Its `add` method always accepts the event: - **Before the first listener**, events go into a pending queue that is handed to the subscription when it arrives. - **While the listener is paused**, the subscription buffers every delivered event and replays them on resume. - **After the listener cancels**, `add` does not throw — the event is silently discarded — and `isClosed` is still `false`. The controller exposes the state — `isPaused` is `true` in the first two cases and `hasListener` is `false` in the last — but nothing forces the producer to look. A socket subscription created with `socket.listen((bytes) => controller.add(parse(bytes)))` keeps pumping regardless, so the pause the consumer asked for only moves the backlog from the operating system into Dart heap memory. ## Fix 1: forward the lifecycle by hand Move the socket subscription into the controller's hooks: 1. `onListen`: `upstream = socket.listen(...)` — nothing is read before somebody wants messages. 2. `onPause`: `upstream.pause()` — the socket stops being read; unread bytes stay in the kernel's receive buffer, and TCP's own window stops the server from sending more. 3. `onResume`: `upstream.resume()`. 4. `onCancel`: `return upstream.cancel();` — the socket subscription is released, and the returned future makes the listener's `cancel()` wait for it. Because `onPause` fires only on the first nested pause and `onResume` only after the buffered events have drained, the source is paused exactly as long as the consumer is. ## Fix 2: let the SDK forward it When the controller only relays another stream, hand-wiring is unnecessary: - **`controller.addStream(source)`** subscribes to `source` and forwards its data and errors. The controller propagates its own pause and resume to that inner subscription, and while it runs, direct `add`, `addError`, `close` or another `addStream` throw `StateError('Cannot add event while adding a stream')`. The returned future completes when `source` is done. - **`source.transform(transformer)`**, `map`, `where` and `asyncMap` all return streams whose subscriptions forward pause and cancel upstream, so often the right answer is not to use a controller at all. ## What does not fix it | Attempt | Why it fails | |---|---| | Switching to `StreamController.broadcast()` | Each broadcast subscription still buffers while paused, and events with no listener are simply lost instead | | `sync: true` | Changes when events are delivered, not whether they are buffered during a pause | | Checking `hasListener` only | Stops work after cancel but ignores the paused case | | A bigger device | Delays the crash; the queue is still unbounded | ## How to confirm it - Log `controller.isPaused` inside the socket callback; if it is ever `true` while you are still adding, the producer is ignoring backpressure. - In Flutter DevTools' Memory view, a steadily growing count of your message objects, retained through the stream subscription's pending events, points to the same place. Why a bounded buffer or a drop policy is sometimes preferable to pausing a source is general stream-flow-control theory; the Dart-specific lesson is that **`StreamController` does not apply any of it for you** — the producer must honour `isPaused` or be wired so the SDK forwards it.

  • What happens if code calls add() while controller.addStream() is still running?
    It throws `StateError('Cannot add event while adding a stream')`. Until the future returned by `addStream` completes, the controller accepts events only from that source; `add`, `addError`, `close` and a second `addStream` are all rejected.
  • After the only listener cancels, what does controller.add() do on a default StreamController?
    Nothing visible: the event is discarded without an exception, and `isClosed` stays `false` because the controller was never closed. A producer that keeps working after cancel wastes CPU and I/O, so stop it in `onCancel` or check `hasListener`.

saying these in an interview costs you the question

  • StreamController.add blocks the producer while the listener is paused.
  • Pausing a subscription automatically pauses whatever feeds the controller.
  • Switching to a broadcast controller removes the buffering.
  • add() after the listener cancels throws, so a runaway producer would crash visibly.
  • sync: true stops the controller from buffering during a pause.