Skip to content

Commit

Permalink
worker lock release on terminate
Browse files Browse the repository at this point in the history
  • Loading branch information
nekufa committed Sep 23, 2023
1 parent 25ca5bb commit 89c8fe5
Showing 1 changed file with 8 additions and 2 deletions.
10 changes: 8 additions & 2 deletions sharded_queue/__init__.py
Original file line number Diff line number Diff line change
@@ -1,9 +1,10 @@
from asyncio import sleep
from asyncio import ensure_future, get_event_loop, sleep
from contextlib import asynccontextmanager
from dataclasses import dataclass
from datetime import datetime, timedelta
from functools import cache
from functools import cache, partial
from importlib import import_module
from signal import SIGTERM
from typing import (Any, AsyncGenerator, Generic, NamedTuple, Optional, Self,
TypeVar, get_type_hints)

Expand Down Expand Up @@ -174,10 +175,15 @@ async def loop(
limit: Optional[int] = None,
handler: Optional[type[Handler]] = None,
) -> None:
loop = get_event_loop()
processed = 0
while True and limit is None or limit > processed:
tube = await self.acquire_tube(handler)
loop.add_signal_handler(SIGTERM, partial(
ensure_future, partial(self.lock.release, tube.pipe)
))
processed = processed + await self.process(tube, limit)
loop.remove_signal_handler(SIGTERM)

async def process(self, tube: Tube, limit: Optional[int] = None) -> int:
deserialize = self.queue.serializer.deserialize
Expand Down

0 comments on commit 89c8fe5

Please sign in to comment.