mirror of
https://github.com/django-q2/django-q2.git
synced 2026-09-17 06:27:53 +08:00
Compare commits
5 Commits
post-spawn
...
v1.5.4
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f8501bfb4c | ||
|
|
23c06a38eb | ||
|
|
f8321e7f54 | ||
|
|
630b3c62b3 | ||
|
|
227caa7446 |
14
CHANGELOG.md
14
CHANGELOG.md
@@ -2,6 +2,20 @@
|
|||||||
|
|
||||||
## [Unreleased](https://github.com/GDay/django-q2/tree/HEAD)
|
## [Unreleased](https://github.com/GDay/django-q2/tree/HEAD)
|
||||||
|
|
||||||
|
## [v1.5.4](https://github.com/GDay/django-q2/tree/v1.5.4) (2023-06-29)
|
||||||
|
|
||||||
|
**Merged pull requests:**
|
||||||
|
|
||||||
|
- Rerun successful tasks https://github.com/django-q2/django-q2/pull/99
|
||||||
|
|
||||||
|
## [v1.5.3](https://github.com/GDay/django-q2/tree/v1.5.3) (2023-05-14)
|
||||||
|
|
||||||
|
**Merged pull requests:**
|
||||||
|
|
||||||
|
- Add post_spawn signal. https://github.com/django-q2/django-q2/pull/93
|
||||||
|
- Post spawn docs https://github.com/django-q2/django-q2/pull/95
|
||||||
|
- Make processes identifiable with uuid4 https://github.com/django-q2/django-q2/pull/91
|
||||||
|
|
||||||
## [v1.5.2](https://github.com/GDay/django-q2/tree/v1.5.2) (2023-04-13)
|
## [v1.5.2](https://github.com/GDay/django-q2/tree/v1.5.2) (2023-04-13)
|
||||||
|
|
||||||
**Merged pull requests:**
|
**Merged pull requests:**
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
import django
|
import django
|
||||||
|
|
||||||
VERSION = (1, 5, 2)
|
VERSION = (1, 5, 4)
|
||||||
|
|
||||||
if django.VERSION < (3, 2):
|
if django.VERSION < (3, 2):
|
||||||
default_app_config = "django_q.apps.DjangoQConfig"
|
default_app_config = "django_q.apps.DjangoQConfig"
|
||||||
|
|||||||
@@ -10,10 +10,29 @@ from django_q.models import Failure, OrmQ, Schedule, Success, Task
|
|||||||
from django_q.tasks import async_task
|
from django_q.tasks import async_task
|
||||||
|
|
||||||
|
|
||||||
|
def resubmit_task(model_admin, request, queryset):
|
||||||
|
"""Submit selected tasks back to the queue."""
|
||||||
|
for task in queryset:
|
||||||
|
async_task(
|
||||||
|
task.func,
|
||||||
|
*task.args or (),
|
||||||
|
hook=task.hook,
|
||||||
|
group=task.group,
|
||||||
|
cluster=task.cluster,
|
||||||
|
**task.kwargs or {},
|
||||||
|
)
|
||||||
|
if isinstance(model_admin, FailAdmin):
|
||||||
|
task.delete()
|
||||||
|
|
||||||
|
|
||||||
|
resubmit_task.short_description = _("Resubmit selected tasks to queue")
|
||||||
|
|
||||||
|
|
||||||
class TaskAdmin(admin.ModelAdmin):
|
class TaskAdmin(admin.ModelAdmin):
|
||||||
"""model admin for success tasks."""
|
"""model admin for success tasks."""
|
||||||
|
|
||||||
list_display = ("name", "group", "func", "cluster", "started", "stopped", "time_taken")
|
list_display = ("name", "group", "func", "cluster", "started", "stopped", "time_taken")
|
||||||
|
actions = [resubmit_task]
|
||||||
|
|
||||||
def has_add_permission(self, request):
|
def has_add_permission(self, request):
|
||||||
"""Don't allow adds."""
|
"""Don't allow adds."""
|
||||||
@@ -33,17 +52,6 @@ class TaskAdmin(admin.ModelAdmin):
|
|||||||
return list(self.readonly_fields) + [field.name for field in obj._meta.fields]
|
return list(self.readonly_fields) + [field.name for field in obj._meta.fields]
|
||||||
|
|
||||||
|
|
||||||
def retry_failed(FailAdmin, request, queryset):
|
|
||||||
"""Submit selected tasks back to the queue."""
|
|
||||||
for task in queryset:
|
|
||||||
async_task(task.func, *task.args or (), hook=task.hook,
|
|
||||||
group=task.group, cluster=task.cluster, **task.kwargs or {})
|
|
||||||
task.delete()
|
|
||||||
|
|
||||||
|
|
||||||
retry_failed.short_description = _("Resubmit selected tasks to queue")
|
|
||||||
|
|
||||||
|
|
||||||
class FailAdmin(admin.ModelAdmin):
|
class FailAdmin(admin.ModelAdmin):
|
||||||
"""model admin for failed tasks."""
|
"""model admin for failed tasks."""
|
||||||
|
|
||||||
@@ -53,7 +61,7 @@ class FailAdmin(admin.ModelAdmin):
|
|||||||
"""Don't allow adds."""
|
"""Don't allow adds."""
|
||||||
return False
|
return False
|
||||||
|
|
||||||
actions = [retry_failed]
|
actions = [resubmit_task]
|
||||||
search_fields = ("name", "func", "group")
|
search_fields = ("name", "func", "group")
|
||||||
list_filter = ("group", "cluster")
|
list_filter = ("group", "cluster")
|
||||||
readonly_fields = []
|
readonly_fields = []
|
||||||
|
|||||||
@@ -69,6 +69,7 @@ class Cluster:
|
|||||||
self.start_event = Event()
|
self.start_event = Event()
|
||||||
self.sentinel = Process(
|
self.sentinel = Process(
|
||||||
target=Sentinel,
|
target=Sentinel,
|
||||||
|
name=f"Process-{uuid.uuid4().hex}",
|
||||||
args=(
|
args=(
|
||||||
self.stop_event,
|
self.stop_event,
|
||||||
self.start_event,
|
self.start_event,
|
||||||
@@ -196,7 +197,7 @@ class Sentinel:
|
|||||||
"""
|
"""
|
||||||
:type target: function or class
|
:type target: function or class
|
||||||
"""
|
"""
|
||||||
p = Process(target=target, args=args)
|
p = Process(target=target, args=args, name=f"Process-{uuid.uuid4().hex}")
|
||||||
p.daemon = True
|
p.daemon = True
|
||||||
if target == worker:
|
if target == worker:
|
||||||
p.daemon = Conf.DAEMONIZE_WORKERS
|
p.daemon = Conf.DAEMONIZE_WORKERS
|
||||||
|
|||||||
@@ -64,7 +64,7 @@ def test_admin_views(admin_client, monkeypatch):
|
|||||||
|
|
||||||
# resubmit the failure
|
# resubmit the failure
|
||||||
url = reverse("admin:django_q_failure_changelist")
|
url = reverse("admin:django_q_failure_changelist")
|
||||||
data = {"action": "retry_failed", "_selected_action": [f.pk]}
|
data = {"action": "resubmit_task", "_selected_action": [f.pk]}
|
||||||
response = admin_client.post(url, data)
|
response = admin_client.post(url, data)
|
||||||
assert response.status_code == 302
|
assert response.status_code == 302
|
||||||
assert Failure.objects.filter(name=f.id).exists() is False
|
assert Failure.objects.filter(name=f.id).exists() is False
|
||||||
@@ -84,3 +84,10 @@ def test_admin_views(admin_client, monkeypatch):
|
|||||||
data = {"post": "yes"}
|
data = {"post": "yes"}
|
||||||
response = admin_client.post(url, data)
|
response = admin_client.post(url, data)
|
||||||
assert response.status_code == 302
|
assert response.status_code == 302
|
||||||
|
# Resubmit a successful task.
|
||||||
|
url = reverse("admin:django_q_success_changelist")
|
||||||
|
data = {"action": "resubmit_task", "_selected_action": [t.pk]}
|
||||||
|
initial_queue_count = OrmQ.objects.count()
|
||||||
|
response = admin_client.post(url, data)
|
||||||
|
assert response.status_code == 302
|
||||||
|
assert OrmQ.objects.count() > initial_queue_count
|
||||||
|
|||||||
@@ -11,36 +11,36 @@ Start your cluster using Django's ``manage.py`` command::
|
|||||||
|
|
||||||
You should see the cluster starting ::
|
You should see the cluster starting ::
|
||||||
|
|
||||||
10:57:40 [Q] INFO Q Cluster-31781 starting.
|
10:57:40 [Q] INFO Q Cluster freddie-uncle-twenty-ten starting.
|
||||||
10:57:40 [Q] INFO Process-1:1 ready for work at 31784
|
10:57:40 [Q] INFO Process-ede257774c4444c980ab479f10947acc ready for work at 31784
|
||||||
10:57:40 [Q] INFO Process-1:2 ready for work at 31785
|
10:57:40 [Q] INFO Process-ed580482da3f42968230baa2e4253e42 ready for work at 31785
|
||||||
10:57:40 [Q] INFO Process-1:3 ready for work at 31786
|
10:57:40 [Q] INFO Process-8a370dc2bc1d49aa9864e517c9895f74 ready for work at 31786
|
||||||
10:57:40 [Q] INFO Process-1:4 ready for work at 31787
|
10:57:40 [Q] INFO Process-74912f9844264d1397c6e54476b530c0 ready for work at 31787
|
||||||
10:57:40 [Q] INFO Process-1:5 ready for work at 31788
|
10:57:40 [Q] INFO Process-b00edb26c6074a6189e5696c60aeb35b ready for work at 31788
|
||||||
10:57:40 [Q] INFO Process-1:6 ready for work at 31789
|
10:57:40 [Q] INFO Process-b0862965db04479f9784a26639ee51e0 ready for work at 31789
|
||||||
10:57:40 [Q] INFO Process-1:7 ready for work at 31790
|
10:57:40 [Q] INFO Process-7e8abbb8ca2d4d9bb20a937dd5e2872b ready for work at 31790
|
||||||
10:57:40 [Q] INFO Process-1:8 ready for work at 31791
|
10:57:40 [Q] INFO Process-b0862965db04479f9784a26639ee51e0 ready for work at 31791
|
||||||
10:57:40 [Q] INFO Process-1:9 monitoring at 31792
|
10:57:40 [Q] INFO Process-67fa9461ac034736a766cd813f617e62 monitoring at 31792
|
||||||
10:57:40 [Q] INFO Process-1 guarding cluster at 31783
|
10:57:40 [Q] INFO Process-eac052c646b2459797cee98bdb84c85d guarding cluster at 31783
|
||||||
10:57:40 [Q] INFO Process-1:10 pushing tasks at 31793
|
10:57:40 [Q] INFO Process-5d98deb19b1e4b2da2ef1e5bd6824f75 pushing tasks at 31793
|
||||||
10:57:40 [Q] INFO Q Cluster-31781 running.
|
10:57:40 [Q] INFO Q Cluster freddie-uncle-twenty-ten running.
|
||||||
|
|
||||||
|
|
||||||
Stopping the cluster with ctrl-c or either the ``SIGTERM`` and ``SIGKILL`` signals, will initiate the :ref:`stop_procedure`::
|
Stopping the cluster with ctrl-c or either the ``SIGTERM`` and ``SIGKILL`` signals, will initiate the :ref:`stop_procedure`::
|
||||||
|
|
||||||
16:44:12 [Q] INFO Q Cluster-31781 stopping.
|
16:44:12 [Q] INFO Q Cluster freddie-uncle-twenty-ten stopping.
|
||||||
16:44:12 [Q] INFO Process-1 stopping cluster processes
|
16:44:12 [Q] INFO Process-eac052c646b2459797cee98bdb84c85d stopping cluster processes
|
||||||
16:44:13 [Q] INFO Process-1:10 stopped pushing tasks
|
16:44:13 [Q] INFO Process-5d98deb19b1e4b2da2ef1e5bd6824f75 stopped pushing tasks
|
||||||
16:44:13 [Q] INFO Process-1:6 stopped doing work
|
16:44:13 [Q] INFO Process-b0862965db04479f9784a26639ee51e0 stopped doing work
|
||||||
16:44:13 [Q] INFO Process-1:4 stopped doing work
|
16:44:13 [Q] INFO Process-7e8abbb8ca2d4d9bb20a937dd5e2872b stopped doing work
|
||||||
16:44:13 [Q] INFO Process-1:1 stopped doing work
|
16:44:13 [Q] INFO Process-b0862965db04479f9784a26639ee51e0 stopped doing work
|
||||||
16:44:13 [Q] INFO Process-1:5 stopped doing work
|
16:44:13 [Q] INFO Process-b00edb26c6074a6189e5696c60aeb35b stopped doing work
|
||||||
16:44:13 [Q] INFO Process-1:7 stopped doing work
|
16:44:13 [Q] INFO Process-74912f9844264d1397c6e54476b530c0 stopped doing work
|
||||||
16:44:13 [Q] INFO Process-1:3 stopped doing work
|
16:44:13 [Q] INFO Process-8a370dc2bc1d49aa9864e517c9895f74 stopped doing work
|
||||||
16:44:13 [Q] INFO Process-1:8 stopped doing work
|
16:44:13 [Q] INFO Process-ed580482da3f42968230baa2e4253e42 stopped doing work
|
||||||
16:44:13 [Q] INFO Process-1:2 stopped doing work
|
16:44:13 [Q] INFO Process-ede257774c4444c980ab479f10947acc stopped doing work
|
||||||
16:44:14 [Q] INFO Process-1:9 stopped monitoring results
|
16:44:14 [Q] INFO Process-67fa9461ac034736a766cd813f617e62 stopped monitoring results
|
||||||
16:44:15 [Q] INFO Q Cluster-31781 has stopped.
|
16:44:15 [Q] INFO Q Cluster freddie-uncle-twenty-ten has stopped.
|
||||||
|
|
||||||
The number of workers, optional timeouts, recycles and cpu_affinity can be controlled via the :doc:`configure` settings.
|
The number of workers, optional timeouts, recycles and cpu_affinity can be controlled via the :doc:`configure` settings.
|
||||||
|
|
||||||
|
|||||||
@@ -75,7 +75,7 @@ author = "Ilan Steemers, Stan Triepels"
|
|||||||
# The short X.Y version.
|
# The short X.Y version.
|
||||||
version = "1.5"
|
version = "1.5"
|
||||||
# The full version, including alpha/beta/rc tags.
|
# The full version, including alpha/beta/rc tags.
|
||||||
release = "1.5.2"
|
release = "1.5.4"
|
||||||
|
|
||||||
# The language for content autogenerated by Sphinx. Refer to documentation
|
# The language for content autogenerated by Sphinx. Refer to documentation
|
||||||
# for a list of supported languages.
|
# for a list of supported languages.
|
||||||
|
|||||||
@@ -13,6 +13,12 @@ Before enqueuing a task
|
|||||||
The ``django_q.signals.pre_enqueue`` signal is emitted before a task is
|
The ``django_q.signals.pre_enqueue`` signal is emitted before a task is
|
||||||
enqueued. The task dictionary is given as the ``task`` argument.
|
enqueued. The task dictionary is given as the ``task`` argument.
|
||||||
|
|
||||||
|
After spawning a worker process
|
||||||
|
"""""""""""""""""""""""""""""""
|
||||||
|
|
||||||
|
The ``django_q.signals.post_spawn`` signal is emitted after a worker process has
|
||||||
|
spawned. The process name is given as the ``proc_name`` argument (string).
|
||||||
|
|
||||||
Before executing a task
|
Before executing a task
|
||||||
"""""""""""""""""""""""
|
"""""""""""""""""""""""
|
||||||
|
|
||||||
@@ -37,7 +43,7 @@ Connecting to a Django Q2 signal is done the same as any other Django
|
|||||||
signal::
|
signal::
|
||||||
|
|
||||||
from django.dispatch import receiver
|
from django.dispatch import receiver
|
||||||
from django_q.signals import pre_enqueue, pre_execute, post_execute
|
from django_q.signals import pre_enqueue, pre_execute, post_execute, post_spawn
|
||||||
|
|
||||||
@receiver(pre_enqueue)
|
@receiver(pre_enqueue)
|
||||||
def my_pre_enqueue_callback(sender, task, **kwargs):
|
def my_pre_enqueue_callback(sender, task, **kwargs):
|
||||||
@@ -51,4 +57,8 @@ signal::
|
|||||||
def my_post_execute_callback(sender, task, **kwargs):
|
def my_post_execute_callback(sender, task, **kwargs):
|
||||||
print(f"Task {task['name']} was executed with result {task['result']}")
|
print(f"Task {task['name']} was executed with result {task['result']}")
|
||||||
|
|
||||||
|
@receiver(post_spawn)
|
||||||
|
def my_post_spawn_callback(sender, proc_name, **kwargs):
|
||||||
|
print(f"Process {proc_name} has spawned")
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[tool.poetry]
|
[tool.poetry]
|
||||||
name = "django-q2"
|
name = "django-q2"
|
||||||
version = "1.5.2"
|
version = "1.5.4"
|
||||||
packages = [
|
packages = [
|
||||||
{ include = "django_q" },
|
{ include = "django_q" },
|
||||||
]
|
]
|
||||||
|
|||||||
Reference in New Issue
Block a user