Skip to content

Commit 99f643d

Browse files
committed
fix error with run_after enqueueing
1 parent f5e19c2 commit 99f643d

5 files changed

Lines changed: 36 additions & 33 deletions

File tree

README.md

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -4,8 +4,7 @@ Steady Queue is a port to Django of the excellent [Solid
44
Queue][solid-queue-github] DB-based queueing backend for Ruby on Rails. The goal
55
of this port has been to keep the internals as close as possible to a direct
66
translation from Ruby to Python, while adapting the external interfaces to be
7-
idiomatic in Django and Python. Read more on the differences between Steady
8-
Queue and Solid Queue below.
7+
idiomatic in Django. Read more on the [differences between Steady Queue and Solid Queue below](#deviations-from-solid-queue).
98

109
Steady queue exposes a task backend that is compatible with the background worker specification outlined in [DEP 0014][DEP0014]. In addition to task enqueueing and processing, it supports delayed tasks, concurrency controls, recurring tasks, pausing queues, numeric priorities per task, priorities by queue order and bulk enqueueing.
1110

@@ -39,7 +38,7 @@ Steady Queue works like any other DEP 0014-compatible task backend.
3938
Tasks are functions decorated with the `@task` decorator from `django.tasks`:
4039

4140
```python
42-
from django.task import task
41+
from django.tasks import task
4342

4443
@task()
4544
def greet(name: str, times: int = 1):

steady_queue/backend.py

Lines changed: 12 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,5 @@
1+
from typing import Any
2+
13
from django.tasks import Task, TaskResult, TaskResultStatus
24
from django.tasks.backends.base import BaseTaskBackend
35

@@ -17,23 +19,25 @@ def validate_task(self, task: Task) -> None:
1719
# TODO: do we need to do anything here?
1820
super().validate_task(task)
1921

20-
def enqueue(self, task, args, kwargs) -> TaskResult:
22+
def enqueue(
23+
self, task: SteadyQueueTask, args: list, kwargs: dict[str, Any]
24+
) -> TaskResult:
2125
from steady_queue.models import Job
2226

2327
if not isinstance(task, SteadyQueueTask):
2428
raise ValueError("Steady Queue only supports SteadyQueueTasks")
2529

26-
task.args = args
27-
task.kwargs = kwargs
28-
job = Job.objects.enqueue(task, scheduled_at=task.run_after)
29-
return self._to_task_result(task, job)
30+
job = Job.objects.enqueue(task, args, kwargs)
31+
return self._to_task_result(task, job, args, kwargs)
3032

3133
def get_result(self, result_id: str) -> TaskResult:
3234
raise NotImplementedError(
3335
"This backend does not support retrieving or refreshing results."
3436
)
3537

36-
def _to_task_result(self, task: SteadyQueueTask, job) -> TaskResult:
38+
def _to_task_result(
39+
self, task: SteadyQueueTask, job, args: list, kwargs: dict[str, Any]
40+
) -> TaskResult:
3741
return TaskResult(
3842
task=task,
3943
id=str(job.id),
@@ -42,8 +46,8 @@ def _to_task_result(self, task: SteadyQueueTask, job) -> TaskResult:
4246
started_at=None,
4347
finished_at=job.finished_at,
4448
last_attempted_at=None,
45-
args=[],
46-
kwargs=task.arguments,
49+
args=args,
50+
kwargs=kwargs,
4751
backend=task.backend,
4852
errors=[],
4953
worker_ids=[],

steady_queue/models/job.py

Lines changed: 8 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,3 @@
1-
from datetime import datetime
21
from typing import Optional
32

43
from django.db import models
@@ -10,11 +9,11 @@
109

1110

1211
class JobQuerySet(ExecutableQuerySet, models.QuerySet):
13-
def enqueue(self, task: SteadyQueueTask, scheduled_at: Optional[datetime] = None):
12+
def enqueue(self, task: SteadyQueueTask, args: list, kwargs: dict):
1413
try:
15-
enqueued_job = self.create(**self.model.attributes_from_django_task(task))
16-
task.provider_task_id = enqueued_job.id
17-
return enqueued_job
14+
return self.create(
15+
**self.model.attributes_from_django_task(task, args, kwargs)
16+
)
1817
except Exception as e:
1918
# TODO: enqueue error
2019
raise e
@@ -64,14 +63,16 @@ class Meta:
6463
DEFAULT_PRIORITY = 0
6564

6665
@classmethod
67-
def attributes_from_django_task(cls, task: SteadyQueueTask):
66+
def attributes_from_django_task(
67+
cls, task: SteadyQueueTask, args: list, kwargs: dict
68+
):
6869
return {
6970
"queue_name": task.queue_name or cls.DEFAULT_QUEUE_NAME,
7071
"django_task_id": task.id,
7172
"priority": task.priority or cls.DEFAULT_PRIORITY,
7273
"scheduled_at": task.run_after or timezone.now(),
7374
"class_name": task.module_path,
74-
"arguments": task.serialize(),
75+
"arguments": task.serialize(args, kwargs),
7576
"concurrency_key": task.concurrency_key,
7677
}
7778

steady_queue/recurring_task.py

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,4 @@
11
import logging
2-
from dataclasses import replace
32
from typing import Optional
43

54
from steady_queue.configuration import Configuration
@@ -41,13 +40,13 @@ def wrapper(task: SteadyQueueTask):
4140
key=key,
4241
class_name=class_name,
4342
schedule=schedule,
44-
arguments=task.serialize(),
43+
arguments=task.serialize(args, kwargs),
4544
queue_name=queue_name,
4645
priority=priority,
4746
description=description,
4847
)
4948
configurations.append(configuration)
5049

51-
return replace(task, args=args, kwargs=kwargs)
50+
return task
5251

5352
return wrapper

steady_queue/task.py

Lines changed: 12 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import datetime
22
import uuid
3-
from dataclasses import dataclass, replace
3+
from dataclasses import dataclass
44
from typing import Any, Optional
55

66
from django.tasks import Task, TaskResult
@@ -16,9 +16,6 @@ class UnknownTaskClassError(Exception):
1616

1717
@dataclass(frozen=True, slots=True, kw_only=True)
1818
class SteadyQueueTask(Task):
19-
args: Optional[list[Any]] = None
20-
kwargs: Optional[dict[str, Any]] = None
21-
2219
concurrency_key: Optional[str] = None
2320
concurrency_limit: Optional[int] = None
2421
concurrency_duration: Optional[timezone.timedelta] = None
@@ -42,7 +39,7 @@ def using(
4239
if isinstance(run_after, datetime.timedelta):
4340
run_after = timezone.now() + run_after
4441

45-
return super().using(
42+
return super(SteadyQueueTask, self).using(
4643
priority=priority,
4744
queue_name=queue_name,
4845
run_after=run_after,
@@ -52,14 +49,14 @@ def using(
5249
def enqueue(self, *args: Any, **kwargs: Any) -> TaskResult:
5350
return self.get_backend().enqueue(self, args, kwargs)
5451

55-
def serialize(self):
52+
def serialize(self, args: list, kwargs: dict):
5653
return {
5754
"class_name": self.module_path,
5855
"job_id": self.id, # TODO: make it stable
5956
"backend": self.backend,
6057
"queue_name": self.queue_name,
6158
"priority": self.priority,
62-
"arguments": Arguments.serialize_args_and_kwargs(self.args, self.kwargs),
59+
"arguments": Arguments.serialize_args_and_kwargs(args, kwargs),
6360
"locale": translation.get_language(),
6461
"timezone": timezone.get_current_timezone_name(),
6562
"enqueued_at": timezone.now().isoformat(),
@@ -68,8 +65,9 @@ def serialize(self):
6865

6966
@classmethod
7067
def execute(cls, job_data: dict[str, Any]):
68+
args, kwargs = Arguments.deserialize_args_and_kwargs(job_data["arguments"])
7169
task = cls.deserialize(job_data)
72-
task.func(*task.args, **task.kwargs)
70+
task.func(*args, **kwargs)
7371

7472
@classmethod
7573
def deserialize(cls, job_data: dict[str, Any]):
@@ -78,11 +76,13 @@ def deserialize(cls, job_data: dict[str, Any]):
7876
except ImportError as e:
7977
raise UnknownTaskClassError(job_data["class_name"]) from e
8078

81-
task = task_class.using(
79+
return task_class.using(
8280
priority=job_data["priority"],
8381
queue_name=job_data["queue_name"],
84-
run_after=job_data["scheduled_at"],
82+
run_after=(
83+
timezone.datetime.fromisoformat(job_data["scheduled_at"])
84+
if job_data["scheduled_at"]
85+
else None
86+
),
8587
backend=job_data["backend"],
8688
)
87-
args, kwargs = Arguments.deserialize_args_and_kwargs(job_data["arguments"])
88-
return replace(task, args=args, kwargs=kwargs)

0 commit comments

Comments
 (0)