Repository navigation
Bound queues and cancel reader tasks in streams.concat - #172
princeraj2572 wants to merge 1 commit into
Conversation
Add an optional queue_maxsize to concat so a full buffer applies backpressure to its input stream instead of buffering everything. Cancel and await the background reader tasks when the consumer exits early, and send the end-of-stream marker without blocking on a full queue during cancellation. Addresses google-gemini#164.
|
Thanks for your pull request! It looks like this may be your first contribution to a Google open source project. Before we can look at your pull request, you'll need to sign a Contributor License Agreement (CLA). View this failed invocation of the CLA check for more information. For the most up to date status, view the checks section at the bottom of the pull request. |
There was a problem hiding this comment.
Code Review
This pull request introduces support for bounded queues in the concat stream processor, allowing backpressure to be applied, and improves task cleanup on early exit. However, the current implementation of task cleanup silently swallows exceptions raised by the input streams because asyncio.gather is called with return_exceptions=True. It is recommended to use asyncio.TaskGroup to ensure proper exception propagation and task cleanup, and to add a corresponding test case to verify exception handling.
| tasks = [ | ||
| context.create_task(_stream_outputs(c, q)) | ||
| for c, q in zip(contents, output_queues) | ||
| ] | ||
|
|
||
| for q in output_queues: | ||
| while (part := await q.get()) is not None: | ||
| q.task_done() | ||
| yield part | ||
| try: | ||
| for q in output_queues: | ||
| async for part in dequeue(q): | ||
| yield part | ||
| finally: | ||
| for t in tasks: | ||
| t.cancel() | ||
| await asyncio.gather(*tasks, return_exceptions=True) |
There was a problem hiding this comment.
The current implementation manually manages background tasks and cleans them up using asyncio.gather(*tasks, return_exceptions=True). However, this has a critical correctness issue: any exceptions raised by the input streams (e.g., ValueError, ConnectionError) are silently swallowed.
When a reader task fails, it catches the exception, puts None into the queue (if not full), and re-raises the exception. The consumer reads None, assumes the stream finished successfully, and moves on. In the finally block, asyncio.gather is called with return_exceptions=True, which returns the exception instead of raising it, and the returned exception is discarded. This results in silent data truncation/loss without any error propagation to the caller.
Using asyncio.TaskGroup (similar to how it is used in merge) elegantly solves this issue. It automatically propagates any exceptions raised by the background tasks to the consumer, and ensures all tasks are cancelled and cleaned up if the consumer exits early or is cancelled.
async with asyncio.TaskGroup() as tg:
for c, q in zip(contents, output_queues):
tg.create_task(_stream_outputs(c, q))
for q in output_queues:
async for part in dequeue(q):
yield part|
|
||
| leaked = [t for t in asyncio.all_tasks() - before if not t.done()] | ||
| self.assertEmpty(leaked) | ||
| self.assertCountEqual(cancelled, ['a', 'b']) |
There was a problem hiding this comment.
Please add a test case to verify that exceptions raised by the input streams are correctly propagated to the consumer instead of being silently swallowed.
self.assertCountEqual(cancelled, ['a', 'b'])
async def test_concat_propagates_exceptions(self):
async def _error_producer():
yield content_api.ProcessorPart('1')
raise ValueError('producer error')
with self.assertRaises(ExceptionGroup) as ctx:
await streams.gather_stream(
streams.concat(_error_producer())
)
self.assertIsInstance(ctx.exception.exceptions[0], ValueError)
Addresses #164.
Problem
streams.concatread every input stream into an unboundedasyncio.Queueand never cancelled its reader tasks. A slow consumer let memory grow with the whole output of the later streams, and a caller that stopped early left the reader tasks running.Changes
queue_maxsizeargument (default0, unbounded, so existing behaviour is unchanged). When set, a full buffer pauses its input stream until the consumer catches up.Tests
Added tests for backpressure with a bounded queue, ordering with a small bound, and no leaked tasks after breaking out early.
processor_test.pyandstreams_test.pypass.