From 856112b1c2ebe6640e76c51d5855d926f090c521 Mon Sep 17 00:00:00 2001 From: dmitry krokhin Date: Mon, 25 Sep 2023 12:50:29 +0300 Subject: [PATCH] worker loop optional tube acquire --- sharded_queue/__init__.py | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/sharded_queue/__init__.py b/sharded_queue/__init__.py index 63c6b08..39c20a4 100644 --- a/sharded_queue/__init__.py +++ b/sharded_queue/__init__.py @@ -136,7 +136,7 @@ class Worker: async def acquire_tube( self, handler: Optional[type[Handler]] = None - ) -> Tube: + ) -> Optional[Tube]: all_pipes = False while get_event_loop().is_running(): for pipe in await self.queue.storage.pipes(): @@ -165,6 +165,8 @@ async def acquire_tube( else: all_pipes = True + return None + def page_size(self, limit: Optional[int] = None) -> int: if limit is None: return settings.worker_batch_size @@ -182,6 +184,8 @@ async def loop( limit is None or limit > processed ): tube = await self.acquire_tube(handler) + if not tube: + break self.pipe = tube.pipe processed = processed + await self.process(tube, limit) self.pipe = None