(self, task)
| 296 | return self._registry.create_task(message) |
| 297 | |
| 298 | def enqueue(self, task): |
| 299 | if isinstance(task, group): |
| 300 | return ResultGroup([self.enqueue(t) for t in task.tasks]) |
| 301 | elif isinstance(task, chord): |
| 302 | return self._enqueue_chord(task) |
| 303 | |
| 304 | # Resolve the expiration time when the task is enqueued. |
| 305 | if task.expires: |
| 306 | task.resolve_expires(self.utc) |
| 307 | |
| 308 | self._emit(S.SIGNAL_ENQUEUED, task) |
| 309 | |
| 310 | if self._immediate: |
| 311 | self.execute(task) |
| 312 | else: |
| 313 | self.storage.enqueue(self.serialize_task(task), task.priority) |
| 314 | |
| 315 | if not self.results: |
| 316 | return |
| 317 | |
| 318 | if task.on_complete: |
| 319 | current = task |
| 320 | results = [] |
| 321 | while current is not None: |
| 322 | results.append(Result(self, current)) |
| 323 | current = current.on_complete |
| 324 | return ResultGroup(results) |
| 325 | else: |
| 326 | return Result(self, task) |
| 327 | |
| 328 | def _enqueue_chord(self, chord_obj): |
| 329 | cid = str(uuid.uuid4()) |
no test coverage detected