summaryrefslogtreecommitdiff
path: root/upstream-layers/bitbake/lib/bb/asyncrpc/taskgroup.py
blob: b61f4038597e8b28bd8c43f5ad1ff60a22a56149 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
#
# Copyright BitBake Contributors
#
# SPDX-License-Identifier: GPL-2.0-only
#
import asyncio
import logging

logger = logging.getLogger("asyncio.TaskGroup")


class TaskGroup(object):
    def __init__(self):
        self._tasks = []

    def create_task(self, coro, **kwargs):
        self._tasks.append(asyncio.create_task(coro, **kwargs))

    async def __aenter__(self):
        return self

    async def __aexit__(self, exc_type, exc, tb):
        try:
            if exc is None:
                while self._tasks:
                    done, pending = await asyncio.wait(
                        self._tasks, return_when=asyncio.FIRST_COMPLETED
                    )
                    self._tasks = pending

                    for t in done:
                        try:
                            await t
                        except asyncio.CancelledError:
                            pass

        finally:
            for t in self._tasks:
                t.cancel()
                try:
                    await t
                except:
                    # Ignore exceptions
                    pass

        return False

    @classmethod
    async def run(cls, *coros):
        async with cls() as group:
            for c in coros:
                group.create_task(c)