2019-12-02 23:26:49 +00:00
|
|
|
# Copyright 2019 Canonical, Ltd.
|
|
|
|
#
|
|
|
|
# This program is free software: you can redistribute it and/or modify
|
|
|
|
# it under the terms of the GNU Affero General Public License as
|
|
|
|
# published by the Free Software Foundation, either version 3 of the
|
|
|
|
# License, or (at your option) any later version.
|
|
|
|
#
|
|
|
|
# This program is distributed in the hope that it will be useful,
|
|
|
|
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
|
|
|
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
|
|
|
# GNU Affero General Public License for more details.
|
|
|
|
#
|
|
|
|
# You should have received a copy of the GNU Affero General Public License
|
|
|
|
# along with this program. If not, see <http://www.gnu.org/licenses/>.
|
|
|
|
|
|
|
|
import asyncio
|
2020-04-01 23:08:41 +00:00
|
|
|
import concurrent.futures
|
2019-12-15 23:13:37 +00:00
|
|
|
import logging
|
|
|
|
|
|
|
|
|
|
|
|
log = logging.getLogger("subiquitycore.async_helpers")
|
2019-12-02 23:26:49 +00:00
|
|
|
|
|
|
|
|
2019-12-11 02:24:48 +00:00
|
|
|
def _done(fut):
|
|
|
|
try:
|
|
|
|
fut.result()
|
|
|
|
except asyncio.CancelledError:
|
|
|
|
pass
|
|
|
|
|
|
|
|
|
2019-12-14 20:16:13 +00:00
|
|
|
def schedule_task(coro, propagate_errors=True):
|
2022-10-07 16:13:30 +00:00
|
|
|
loop = asyncio.get_running_loop()
|
2019-12-11 02:24:48 +00:00
|
|
|
if asyncio.iscoroutine(coro):
|
|
|
|
task = asyncio.Task(coro)
|
|
|
|
else:
|
|
|
|
task = coro
|
2019-12-14 20:16:13 +00:00
|
|
|
if propagate_errors:
|
|
|
|
task.add_done_callback(_done)
|
2019-12-11 02:24:48 +00:00
|
|
|
loop.call_soon(asyncio.ensure_future, task)
|
|
|
|
return task
|
2019-12-02 23:26:49 +00:00
|
|
|
|
|
|
|
|
2023-01-19 11:30:33 +00:00
|
|
|
# Collection of tasks that we want to fire and forget.
|
|
|
|
# Keeping a reference to all background tasks ensures that the tasks don't get
|
|
|
|
# garbage collected before they are done.
|
|
|
|
# https://docs.python.org/3/library/asyncio-task.html#asyncio.create_task
|
|
|
|
background_tasks = set()
|
|
|
|
|
|
|
|
|
|
|
|
def run_bg_task(coro, *args, **kwargs) -> None:
|
|
|
|
""" Run a background task in a fire-and-forget style. """
|
|
|
|
task = asyncio.create_task(coro, *args, **kwargs)
|
|
|
|
background_tasks.add(task)
|
|
|
|
task.add_done_callback(background_tasks.discard)
|
|
|
|
|
|
|
|
|
2019-12-02 23:26:49 +00:00
|
|
|
async def run_in_thread(func, *args):
|
2022-10-07 16:13:30 +00:00
|
|
|
loop = asyncio.get_running_loop()
|
2020-04-01 23:08:41 +00:00
|
|
|
try:
|
|
|
|
return await loop.run_in_executor(None, func, *args)
|
|
|
|
except concurrent.futures.CancelledError:
|
|
|
|
raise asyncio.CancelledError
|
2019-12-11 02:24:48 +00:00
|
|
|
|
|
|
|
|
2022-05-24 19:31:51 +00:00
|
|
|
class TaskAlreadyRunningError(Exception):
|
|
|
|
"""Used to let callers know that a task hasn't been started due to
|
|
|
|
cancel_restart == False and the task already running."""
|
|
|
|
pass
|
|
|
|
|
|
|
|
|
2019-12-11 02:24:48 +00:00
|
|
|
class SingleInstanceTask:
|
|
|
|
|
2022-05-24 19:31:51 +00:00
|
|
|
def __init__(self, func, propagate_errors=True, cancel_restart=True):
|
2019-12-11 02:24:48 +00:00
|
|
|
self.func = func
|
2019-12-14 20:16:13 +00:00
|
|
|
self.propagate_errors = propagate_errors
|
2019-12-11 02:24:48 +00:00
|
|
|
self.task = None
|
2022-05-24 19:31:51 +00:00
|
|
|
# if True, allow subsequent start calls to cancel a running task
|
|
|
|
# raises TaskAlreadyRunningError if we skip starting the task.
|
|
|
|
self.cancel_restart = cancel_restart
|
2019-12-11 02:24:48 +00:00
|
|
|
|
2019-12-15 23:13:37 +00:00
|
|
|
async def _start(self, old):
|
|
|
|
if old is not None:
|
|
|
|
old.cancel()
|
2019-12-11 02:24:48 +00:00
|
|
|
try:
|
2019-12-15 23:13:37 +00:00
|
|
|
await old
|
2019-12-12 10:01:51 +00:00
|
|
|
except BaseException:
|
2019-12-11 02:24:48 +00:00
|
|
|
pass
|
2019-12-15 23:13:37 +00:00
|
|
|
schedule_task(self.task, self.propagate_errors)
|
|
|
|
|
|
|
|
async def start(self, *args, **kw):
|
|
|
|
await self.start_sync(*args, **kw)
|
2019-12-14 20:16:13 +00:00
|
|
|
return self.task
|
2019-12-11 02:24:48 +00:00
|
|
|
|
|
|
|
def start_sync(self, *args, **kw):
|
2022-05-24 19:31:51 +00:00
|
|
|
if not self.cancel_restart:
|
|
|
|
if self.task is not None and not self.task.done():
|
|
|
|
raise TaskAlreadyRunningError(
|
|
|
|
'Skipping invocation of task - already running')
|
2019-12-15 23:13:37 +00:00
|
|
|
old = self.task
|
|
|
|
coro = self.func(*args, **kw)
|
|
|
|
if asyncio.iscoroutine(coro):
|
|
|
|
self.task = asyncio.Task(coro)
|
|
|
|
else:
|
|
|
|
self.task = coro
|
|
|
|
return schedule_task(self._start(old))
|
2019-12-19 03:22:07 +00:00
|
|
|
|
|
|
|
async def wait(self):
|
|
|
|
while True:
|
|
|
|
try:
|
|
|
|
return await self.task
|
|
|
|
except asyncio.CancelledError:
|
|
|
|
pass
|