• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

Morry98 / fastapi-task-manager / 22852159792

09 Mar 2026 11:52AM UTC coverage: 97.577%. First build
22852159792

Pull #2

github

web-flow
Merge 8228b755b into d04f85da2
Pull Request #2: Improve task management with Redis Streams, leader election and management API

3182 of 3264 new or added lines in 32 files covered. (97.49%)

3342 of 3425 relevant lines covered (97.58%)

4.86 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

97.89
/src/fastapi_task_manager/task_group.py
1
import hashlib
5✔
2
import logging
5✔
3
from collections.abc import Callable
5✔
4

5
from fastapi_task_manager.schema.task import Task
5✔
6

7
logger = logging.getLogger("fastapi.task-manager.group")
5✔
8

9

10
class TaskGroup:
5✔
11
    def __init__(
5✔
12
        self,
13
        name: str,
14
        tags: list[str] | None = None,
15
    ):
16
        self._name = name
5✔
17
        self._tags = tags
5✔
18
        self._tasks: list[Task] = []
5✔
19
        # Registry of callable functions available for dynamic task creation
20
        self._function_registry: dict[str, Callable] = {}
5✔
21

22
    @property
5✔
23
    def name(self) -> str:
5✔
24
        return self._name
5✔
25

26
    @property
5✔
27
    def tags(self):
5✔
28
        return self._tags.copy() if self._tags else []
5✔
29

30
    @property
5✔
31
    def tasks(self) -> list[Task]:
5✔
32
        """Get all tasks in the group."""
33
        return self._tasks.copy()
5✔
34

35
    @property
5✔
36
    def function_registry(self) -> dict[str, Callable]:
5✔
37
        """Get a copy of the function registry."""
38
        return self._function_registry.copy()
5✔
39

40
    def register_function(
5✔
41
        self,
42
        name: str | None = None,
43
    ):
44
        """Decorator to register a function in the registry for dynamic task creation.
45

46
        The function is only stored in the registry — no task is created.
47
        Tasks can be created at runtime via the management API referencing this name.
48
        """
49

50
        def wrapper(func: Callable):
5✔
51
            registry_name = name or func.__name__  # ty: ignore[unresolved-attribute]
5✔
52
            if registry_name in self._function_registry:
5✔
53
                msg = f"Function '{registry_name}' is already registered in group '{self._name}'."
5✔
54
                raise RuntimeError(msg)
5✔
55
            self._function_registry[registry_name] = func
5✔
56
            return func
5✔
57

58
        return wrapper
5✔
59

60
    def add_task(  # noqa: PLR0913
5✔
61
        self,
62
        expr: str | list[str],
63
        kwargs: dict | list[dict] | None = None,
64
        tags: list[str] | None = None,
65
        name: str | None = None,
66
        description: str | None = None,
67
        high_priority: bool = False,
68
        retry_backoff: float | None = None,
69
        retry_backoff_max: float | None = None,
70
        register: bool = True,
71
    ):
72
        """Decorator for creating task.
73

74
        By default, the function is also added to the function registry so it can
75
        be used to create dynamic tasks via API. Set register=False to opt out.
76
        """
77

78
        def wrapper(func: Callable):
5✔
79
            # Register the function in the registry unless explicitly disabled
80
            if register:
5✔
81
                registry_name = name or func.__name__  # ty: ignore[unresolved-attribute]
5✔
82
                if registry_name not in self._function_registry:
5✔
83
                    self._function_registry[registry_name] = func
5✔
84

85
            _tags = self._tags or [] + (tags or [])
5✔
86
            # check that both expr and kwargs are lists of the same length
87
            if isinstance(expr, list) and isinstance(kwargs, list) and len(expr) != len(kwargs):
5✔
88
                msg = "expr and kwargs lists must have the same length."
5✔
89
                raise TypeError(msg)
5✔
90
            if isinstance(expr, list):
5✔
91
                expr_list = expr
5✔
92
            else:
93
                expr_list = [expr for _ in range(len(kwargs) if isinstance(kwargs, list) else 1)]
5✔
94
            kwargs_list = kwargs if isinstance(kwargs, list) else [kwargs or {} for _ in range(len(expr_list))]
5✔
95

96
            for i in range(len(expr_list)):
5✔
97
                # add hash of kwargs to expression to make it unique
98
                # pay attention that python hash is not a stable hash across different runs
99
                # so we use sha256 from hashlib
100
                internal_name = name or func.__name__  # ty: ignore[unresolved-attribute]
5✔
101
                hash_suffix = ""
5✔
102
                if kwargs_list[i]:
5✔
103
                    hash_input = str(sorted(kwargs_list[i].items())).encode()
5✔
104
                    hash_suffix = hashlib.sha256(hash_input).hexdigest()
5✔
105
                expr_hash = hashlib.sha256(expr_list[i].encode()).hexdigest()
5✔
106
                internal_name += f"__{hash_suffix}__{expr_hash}__"
5✔
107

108
                for t in self._tasks:
5✔
109
                    if t.name == internal_name:
5✔
110
                        msg = f"Task with name {internal_name} already exists inside group {self.name}."
5✔
111
                        raise RuntimeError(msg)
5✔
112

113
                task = Task(
5✔
114
                    function=func,
115
                    expression=expr_list[i],
116
                    name=internal_name,
117
                    description=description,
118
                    tags=_tags or None,
119
                    high_priority=high_priority,
120
                    kwargs=kwargs_list[i],
121
                    retry_backoff=retry_backoff,
122
                    retry_backoff_max=retry_backoff_max,
123
                    dynamic=False,
124
                )
125
                self._tasks.append(task)
5✔
126

127
            return func
5✔
128

129
        return wrapper
5✔
130

131
    def add_dynamic_task(  # noqa: PLR0913
5✔
132
        self,
133
        function_name: str,
134
        cron_expression: str,
135
        kwargs: dict | None = None,
136
        name: str | None = None,
137
        description: str | None = None,
138
        high_priority: bool = False,
139
        tags: list[str] | None = None,
140
        retry_backoff: float | None = None,
141
        retry_backoff_max: float | None = None,
142
    ) -> Task:
143
        """Create a dynamic task from a registered function.
144

145
        Returns the created Task object.
146
        Raises RuntimeError if the function is not registered or name is not unique.
147
        """
148
        if function_name not in self._function_registry:
5✔
149
            msg = f"Function '{function_name}' is not registered in group '{self._name}'."
5✔
150
            raise RuntimeError(msg)
5✔
151

152
        func = self._function_registry[function_name]
5✔
153
        # Use provided name or generate from function_name + kwargs/expr hashes
154
        task_name = name or function_name
5✔
155
        if not name:
5✔
156
            # Generate unique name using the same hash pattern as add_task
157
            hash_suffix = ""
5✔
158
            if kwargs:
5✔
NEW
159
                hash_input = str(sorted(kwargs.items())).encode()
×
NEW
160
                hash_suffix = hashlib.sha256(hash_input).hexdigest()
×
161
            expr_hash = hashlib.sha256(cron_expression.encode()).hexdigest()
5✔
162
            task_name += f"__{hash_suffix}__{expr_hash}__"
5✔
163

164
        # Validate uniqueness
165
        for t in self._tasks:
5✔
166
            if t.name == task_name:
5✔
167
                msg = f"Task with name '{task_name}' already exists inside group '{self._name}'."
5✔
168
                raise RuntimeError(msg)
5✔
169

170
        # Merge group tags with task-specific tags
171
        merged_tags = (self._tags or []) + (tags or []) or None
5✔
172

173
        task = Task(
5✔
174
            function=func,
175
            expression=cron_expression,
176
            name=task_name,
177
            description=description,
178
            tags=merged_tags,
179
            high_priority=high_priority,
180
            kwargs=kwargs or {},
181
            retry_backoff=retry_backoff,
182
            retry_backoff_max=retry_backoff_max,
183
            dynamic=True,
184
            function_name=function_name,
185
        )
186
        self._tasks.append(task)
5✔
187
        logger.info("Dynamic task '%s' added to group '%s'", task_name, self._name)
5✔
188
        return task
5✔
189

190
    def remove_dynamic_task(self, task_name: str) -> Task:
5✔
191
        """Remove a dynamic task by name.
192

193
        Only dynamic tasks can be removed. Static tasks (registered via decorator)
194
        cannot be removed at runtime.
195

196
        Returns the removed Task object.
197
        Raises RuntimeError if the task is not found or is not dynamic.
198
        """
199
        for i, t in enumerate(self._tasks):
5✔
200
            if t.name == task_name:
5✔
201
                if not t.dynamic:
5✔
202
                    msg = f"Task '{task_name}' in group '{self._name}' is static and cannot be removed at runtime."
5✔
203
                    raise RuntimeError(msg)
5✔
204
                removed = self._tasks.pop(i)
5✔
205
                logger.info("Dynamic task '%s' removed from group '%s'", task_name, self._name)
5✔
206
                return removed
5✔
207

208
        msg = f"Task '{task_name}' not found in group '{self._name}'."
5✔
209
        raise RuntimeError(msg)
5✔
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc