Skip to content

[Bug]: Prevent unbounded queue memory spikes and task leaks in streams.concat #164

Description

@THRISHAL12345

Description of the bug:

In genai_processors/streams.py, streams.concat(*contents) concatenates multiple asynchronous streams by initializing internal asyncio.Queue() instances and spawning concurrent background tasks (_stream_outputs).
Because these queues are unbounded (maxsize=0) and background tasks populate them non-blockingly via put_nowait(c), background tasks immediately consume every item from all input streams into memory. When processing multimodal data (images, audio frames, video segments, PDF documents), this leads to unbounded RAM consumption and potential Out-Of-Memory (OOM) crashes in production environments. Additionally, if stream iteration stops early or encounters an exception, background tasks are never cancelled.

Actual vs expected behavior:

Actual Behaviour

  • Unbounded Memory Growth: streams.concat creates unbounded queues (asyncio.Queue()) and calls put_nowait(c) in background loops. While the consumer is still reading from the first stream (contents[0]), subsequent streams (contents[1], contents[2], etc.) buffer their entire contents in RAM.
  • Task & Resource Leaks: If a consumer stops iterating early (break or aclose()) or an exception is raised, background tasks running _stream_outputs remain running and are not cancelled.
  • Hanging Cancellation: When an enqueue coroutine is cancelled while blocked on a queue put, its finally block attempts await queue.put(None) into a full queue, causing the cancelled task to hang indefinitely.

Expected Behaviour

  • Bounded Memory & Backpressure: streams.concat should enforce bounded buffer capacities (queue_maxsize: int = 1 by default) and use asynchronous put operations (await queue.put(c) via enqueue). Upstream producers should pause once queue_maxsize items are buffered ahead, preventing memory spikes.
  • Clean Task Lifecycle Management: When streams.concat completes iteration or is closed early, all background enqueue tasks should be cleanly cancelled via a try...finally block.
  • Non-Blocking Cancellation: If an enqueue task is cancelled (asyncio.CancelledError), it should immediately re-raise without attempting to block on await queue.put(None).

Any other information you'd like to share?

No response

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions