mirror of
https://github.com/django-q2/django-q2.git
synced 2026-09-15 21:47:53 +08:00
Compare commits
24 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
65d2548e98 | ||
|
|
cbe51f995c | ||
|
|
0090a6f411 | ||
|
|
a154fd658a | ||
|
|
81f2531e3e | ||
|
|
c464751ae9 | ||
|
|
144536cbe4 | ||
|
|
7e72b68c7f | ||
|
|
404932117e | ||
|
|
ddc1aa2de1 | ||
|
|
351bf66d71 | ||
|
|
ccd760a69f | ||
|
|
fb23b6b825 | ||
|
|
35c911080b | ||
|
|
81058fc1cf | ||
|
|
505eb8d537 | ||
|
|
dcfe5b6650 | ||
|
|
c3b49cca25 | ||
|
|
c1e00b09cb | ||
|
|
8e403d0136 | ||
|
|
337a5782d9 | ||
|
|
2c456cd0aa | ||
|
|
8f54e5d5ae | ||
|
|
689c881178 |
48
.github/workflows/test.yml
vendored
48
.github/workflows/test.yml
vendored
@@ -15,30 +15,42 @@ jobs:
|
||||
run: |
|
||||
pipx install ruff==0.4.10
|
||||
ruff format . --check && ruff check .
|
||||
- name: Check twine
|
||||
run: |
|
||||
python -m pip install twine poetry rstcheck
|
||||
poetry build
|
||||
rstcheck README.rst
|
||||
twine check dist/*
|
||||
|
||||
test:
|
||||
runs-on: ubuntu-latest
|
||||
strategy:
|
||||
matrix:
|
||||
python-version: [ "3.8", "3.9", "3.10", "3.11", "3.12" ]
|
||||
django: [ "4.2", "5.0", "5.1"]
|
||||
python-version:
|
||||
- "3.9"
|
||||
- "3.10"
|
||||
- "3.11"
|
||||
- "3.12"
|
||||
django:
|
||||
- "4.2"
|
||||
- "5.0"
|
||||
- "5.1"
|
||||
- "5.2"
|
||||
exclude:
|
||||
# django 5.1 does not support 3.8 and 3.9
|
||||
- python-version: "3.8"
|
||||
django: "5.1"
|
||||
# django 5.2 does not support 3.9
|
||||
- python-version: "3.9"
|
||||
django: "5.2"
|
||||
# django 5.1 does not support 3.9
|
||||
- python-version: "3.9"
|
||||
django: "5.1"
|
||||
# django 5.0 does not support 3.8 and 3.9
|
||||
- python-version: "3.8"
|
||||
django: "5.0"
|
||||
# django 5.0 does not support 3.9
|
||||
- python-version: "3.9"
|
||||
django: "5.0"
|
||||
django: "5.0"
|
||||
# django 4.2 does not support 3.13
|
||||
- python-version: "3.13"
|
||||
django: "4.2"
|
||||
|
||||
services:
|
||||
disque:
|
||||
image: efrecon/disque:1.0-rc1
|
||||
ports:
|
||||
- '7711:7711/tcp'
|
||||
mongodb:
|
||||
image: mongo
|
||||
ports:
|
||||
@@ -67,8 +79,8 @@ jobs:
|
||||
- name: Install dependencies with Django ${{ matrix.django }}
|
||||
run: |
|
||||
python -m pip install --upgrade pip
|
||||
pip install poetry
|
||||
poetry add "django==${{ matrix.django }}" --python=${{ matrix.python-version }}
|
||||
pip install poetry==1.8.5
|
||||
poetry add "django~=${{ matrix.django }}a1" --python=${{ matrix.python-version }}
|
||||
poetry install -E testing
|
||||
- name: Run Tests
|
||||
run: |
|
||||
@@ -93,11 +105,7 @@ jobs:
|
||||
- name: Upload to coveralls
|
||||
run: |
|
||||
python -m pip install --upgrade pip
|
||||
python -m pip install coveralls flake8 black
|
||||
python -m pip install coveralls
|
||||
coveralls --service=github --finish
|
||||
env:
|
||||
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||
- name: Check flake8/black
|
||||
run: |
|
||||
flake8 .
|
||||
black --check .
|
||||
|
||||
43
CHANGELOG.md
43
CHANGELOG.md
@@ -1,6 +1,47 @@
|
||||
# Changelog
|
||||
|
||||
## [v1.7.0](https://github.com/django-q2/django-q2/tree/v1.7.0) (2024-08-24)
|
||||
## [v1.8.0](https://github.com/django-q2/django-q2/tree/v1.7.6) (2025-04-25)
|
||||
|
||||
- Delete deprecated imports and unwrapping https://github.com/django-q2/django-q2/pull/261
|
||||
- Remove disque from CI https://github.com/django-q2/django-q2/pull/270
|
||||
- Avoid creating task run on Scheduler creation https://github.com/django-q2/django-q2/pull/268
|
||||
- Update tested versions, add python 3.13 support and django 5.2 support. Drop python 3.8 support https://github.com/django-q2/django-q2/pull/271
|
||||
- Move timeout function from ORM broker into class to allow easy customization https://github.com/django-q2/django-q2/pull/274
|
||||
|
||||
## [v1.7.6](https://github.com/django-q2/django-q2/tree/v1.7.6) (2025-01-12)
|
||||
|
||||
- Make poetry version fixed in CI https://github.com/django-q2/django-q2/pull/260
|
||||
|
||||
## [v1.7.5](https://github.com/django-q2/django-q2/tree/v1.7.5) (2025-01-12)
|
||||
|
||||
- perf: avoid checking success tasks when save limit is disabled https://github.com/django-q2/django-q2/pull/255
|
||||
- Fix install path for CHANGELOG.md https://github.com/django-q2/django-q2/pull/258
|
||||
|
||||
## [v1.7.4](https://github.com/django-q2/django-q2/tree/v1.7.4) (2024-11-03)
|
||||
|
||||
- Decrease the MAX_RSS set in test_cluster::test_max_rss https://github.com/django-q2/django-q2/pull/240
|
||||
- Fix BROKER_CLASS monkeypatch in test_brokers https://github.com/django-q2/django-q2/pull/239
|
||||
- Fix 'receive_message_wait_time_seconds' SQS broker management https://github.com/django-q2/django-q2/pull/243
|
||||
|
||||
## [v1.7.3](https://github.com/django-q2/django-q2/tree/v1.7.3) (2024-10-15)
|
||||
|
||||
- Catch missing SEGALRM with AttributeError instead of ValueError https://github.com/django-q2/django-q2/pull/223
|
||||
- Refactor timeout handling to handle AttributeError and ValueError for Windows users https://github.com/django-q2/django-q2/pull/234
|
||||
- Fix type check for args in scheduler.py E721 https://github.com/django-q2/django-q2/pull/233
|
||||
- Only trigger prometheus if configured https://github.com/django-q2/django-q2/pull/231
|
||||
- Fix missing ack_id when finishing task https://github.com/django-q2/django-q2/pull/224
|
||||
|
||||
|
||||
## [v1.7.2](https://github.com/django-q2/django-q2/tree/v1.7.2) (2024-09-09)
|
||||
|
||||
- Fix twine check
|
||||
|
||||
## [v1.7.1](https://github.com/django-q2/django-q2/tree/v1.7.1) (2024-09-08)
|
||||
|
||||
- Fixed date of v1.7.0
|
||||
- Fixed README.rst formatting which is blocking release of latest version
|
||||
|
||||
## [v1.7.0](https://github.com/django-q2/django-q2/tree/v1.7.0) (2024-09-08)
|
||||
|
||||
**Merged pull requests:**
|
||||
|
||||
|
||||
@@ -7,6 +7,10 @@ ENV PYTHONUNBUFFERED 1
|
||||
# Sets the default shell to bash
|
||||
ENV SHELL /bin/bash
|
||||
|
||||
RUN set -ex \
|
||||
&& apt update \
|
||||
&& apt-get install gcc python3-dev --yes
|
||||
|
||||
# Upgrades pip
|
||||
RUN pip install -U pip setuptools
|
||||
|
||||
|
||||
8
Makefile
8
Makefile
@@ -2,16 +2,16 @@ dev:
|
||||
docker compose -f web-docker-compose.yaml up
|
||||
|
||||
test:
|
||||
docker-compose -f test-services-docker-compose.yaml run --rm django-q2 poetry run pytest
|
||||
docker compose -f test-services-docker-compose.yaml run --rm django-q2 poetry run pytest
|
||||
|
||||
shell:
|
||||
docker-compose -f test-services-docker-compose.yaml run --rm django-q2 poetry run python manage.py shell
|
||||
docker compose -f test-services-docker-compose.yaml run --rm django-q2 poetry run python manage.py shell
|
||||
|
||||
makemigrations:
|
||||
docker-compose -f test-services-docker-compose.yaml run --rm django-q2 poetry run python manage.py makemigrations
|
||||
docker compose -f test-services-docker-compose.yaml run --rm django-q2 poetry run python manage.py makemigrations
|
||||
|
||||
migrate:
|
||||
docker-compose -f test-services-docker-compose.yaml run --rm django-q2 poetry run python manage.py migrate
|
||||
docker compose -f test-services-docker-compose.yaml run --rm django-q2 poetry run python manage.py migrate
|
||||
|
||||
createsuperuser:
|
||||
docker compose -f web-docker-compose.yaml run --rm web python manage.py createsuperuser
|
||||
|
||||
14
README.rst
14
README.rst
@@ -1,7 +1,7 @@
|
||||
A multiprocessing distributed task queue for Django
|
||||
---------------------------------------------------
|
||||
|
||||
|image0| |image1| |docs| |downloads|
|
||||
|image0| |image1| |downloads|
|
||||
|
||||
Django Q2 is a fork of Django Q. Big thanks to Ilan Steemers for starting this project. Unfortunately, development has stalled since June 2021. Django Q2 is the new updated version of Django Q, with dependencies updates, docs updates and several bug fixes. Original repository: https://github.com/Koed00/django-q
|
||||
|
||||
@@ -38,7 +38,10 @@ Requirements
|
||||
- `Django <https://www.djangoproject.com>`__ > = 4.2
|
||||
- `Django-picklefield <https://github.com/gintas/django-picklefield>`__
|
||||
|
||||
Tested with: Python 3.8, 3.9, 3.10, 3.11 and 3.12. Works with Django 4.2.X and 5.0.X
|
||||
Tested with:
|
||||
|
||||
* Python 3.9 to 3.13.
|
||||
* Django 4.2 to 5.2.
|
||||
|
||||
Brokers
|
||||
~~~~~~~
|
||||
@@ -198,7 +201,7 @@ Admin page or directly from your code:
|
||||
For more info check the `Schedules <https://django-q2.readthedocs.org/en/latest/schedules.html>`__ documentation.
|
||||
|
||||
Development
|
||||
~~~~~~~
|
||||
~~~~~~~~~~~
|
||||
|
||||
There is an example project that you can use to develop with. Docker (compose) is being used to set everything up.
|
||||
Please note that you will have to restart the django-q container when changes have been made to tasks or django-q.
|
||||
@@ -214,7 +217,6 @@ Create a superuser with:
|
||||
|
||||
make createsuperuser
|
||||
|
||||
|
||||
Testing
|
||||
~~~~~~~
|
||||
|
||||
@@ -246,9 +248,5 @@ Acknowledgements
|
||||
:target: https://github.com/GDay/django-q2/actions?query=workflow%3Atests
|
||||
.. |image1| image:: https://coveralls.io/repos/github/GDay/django-q2/badge.svg?branch=master
|
||||
:target: https://coveralls.io/github/GDay/django-q2?branch=master
|
||||
.. |docs| image:: https://readthedocs.org/projects/docs/badge/?version=latest
|
||||
:alt: Documentation Status
|
||||
:scale: 100
|
||||
:target: https://django-q2.readthedocs.org/
|
||||
.. |downloads| image:: https://img.shields.io/pypi/dm/django-q2
|
||||
:target: https://img.shields.io/pypi/dm/django-q2
|
||||
|
||||
19
containers/localstack/init.sh
Executable file
19
containers/localstack/init.sh
Executable file
@@ -0,0 +1,19 @@
|
||||
#!/bin/bash
|
||||
# Note that this file needs to have the executable bit set for it to work with later localstack implementations.
|
||||
|
||||
export DEFAULT_REGION=us-west-2
|
||||
|
||||
create_sqs() {
|
||||
QUEUE_NAME="$1"
|
||||
TIMEOUT=${2:-60}
|
||||
DL_QUEUE_URL=$(awslocal sqs create-queue --queue-name "dl-$QUEUE_NAME" --query QueueUrl --output text)
|
||||
echo ">>> Created $DL_QUEUE_URL queue!"
|
||||
DL_QUEUE_ARN=$(awslocal sqs get-queue-attributes --queue-url "$DL_QUEUE_URL" --attribute-names QueueArn --query Attributes.QueueArn --output text)
|
||||
awslocal sqs create-queue --queue-name "$QUEUE_NAME" --attributes '{
|
||||
"RedrivePolicy": "{\"deadLetterTargetArn\": \"'"$DL_QUEUE_ARN"'\",\"maxReceiveCount\":\"3\"}",
|
||||
"VisibilityTimeout": "'"$TIMEOUT"'"
|
||||
}'
|
||||
}
|
||||
|
||||
# Create SQS queues
|
||||
create_sqs testing
|
||||
@@ -1,8 +1,3 @@
|
||||
import django
|
||||
|
||||
VERSION = (1, 7, 0)
|
||||
|
||||
if django.VERSION < (3, 2):
|
||||
default_app_config = "django_q.apps.DjangoQConfig"
|
||||
VERSION = (1, 8, 0)
|
||||
|
||||
__all__ = ["conf", "cluster", "models", "tasks"]
|
||||
|
||||
@@ -138,7 +138,7 @@ class ScheduleAdmin(admin.ModelAdmin):
|
||||
url = reverse("admin:django_q_success_change", args=(obj.task_id,))
|
||||
else:
|
||||
url = reverse("admin:django_q_failure_change", args=(obj.task_id,))
|
||||
return format_html(f'<a href="{url}">[{obj.task_name}]</a>')
|
||||
return format_html('<a href="{}">[{}]</a>', url, obj.task_name)
|
||||
return None
|
||||
|
||||
get_last_run.allow_tags = True
|
||||
|
||||
@@ -1,3 +1,5 @@
|
||||
import copy
|
||||
|
||||
from boto3 import Session
|
||||
from botocore.client import ClientError
|
||||
|
||||
@@ -78,15 +80,15 @@ class Sqs(Broker):
|
||||
|
||||
@staticmethod
|
||||
def get_connection(list_key: str = None) -> Session:
|
||||
config = Conf.SQS
|
||||
if "aws_region" in config:
|
||||
config["region_name"] = config["aws_region"]
|
||||
del config["aws_region"]
|
||||
config_cloned = copy.deepcopy(Conf.SQS)
|
||||
if "aws_region" in config_cloned:
|
||||
config_cloned["region_name"] = config_cloned["aws_region"]
|
||||
del config_cloned["aws_region"]
|
||||
|
||||
if "receive_message_wait_time_seconds" in config:
|
||||
del config["receive_message_wait_time_seconds"]
|
||||
if "receive_message_wait_time_seconds" in config_cloned:
|
||||
del config_cloned["receive_message_wait_time_seconds"]
|
||||
|
||||
return Session(**config)
|
||||
return Session(**config_cloned)
|
||||
|
||||
def get_queue(self):
|
||||
self.sqs = self.connection.resource("sqs")
|
||||
|
||||
@@ -10,10 +10,6 @@ from django_q.conf import Conf, logger
|
||||
from django_q.models import OrmQ
|
||||
|
||||
|
||||
def _timeout():
|
||||
return timezone.now() + timedelta(seconds=Conf.RETRY)
|
||||
|
||||
|
||||
class ORM(Broker):
|
||||
@staticmethod
|
||||
def get_connection(list_key: str = None):
|
||||
@@ -28,6 +24,9 @@ class ORM(Broker):
|
||||
logger.debug("Broker in an atomic transaction")
|
||||
return OrmQ.objects.using(Conf.ORM)
|
||||
|
||||
def timeout(self, task):
|
||||
return timezone.now() + timedelta(seconds=Conf.RETRY)
|
||||
|
||||
def queue_size(self) -> int:
|
||||
return (
|
||||
self.get_connection()
|
||||
@@ -75,7 +74,7 @@ class ORM(Broker):
|
||||
if (
|
||||
self.get_connection()
|
||||
.filter(id=task.id, lock=task.lock)
|
||||
.update(lock=_timeout())
|
||||
.update(lock=self.timeout(task))
|
||||
):
|
||||
task_list.append((task.pk, task.payload))
|
||||
# else don't process, as another cluster has been faster than us on
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
# Standard
|
||||
import os
|
||||
import signal
|
||||
import socket
|
||||
import uuid
|
||||
@@ -230,8 +231,14 @@ class Sentinel:
|
||||
% {"name": process.name}
|
||||
)
|
||||
else:
|
||||
if prometheus_multiprocess:
|
||||
# check if prometheus is proper configurated
|
||||
prometheus_path = os.getenv(
|
||||
"PROMETHEUS_MULTIPROC_DIR", os.getenv("prometheus_multiproc_dir")
|
||||
)
|
||||
|
||||
if prometheus_multiprocess and prometheus_path:
|
||||
prometheus_multiprocess.mark_process_dead(process.pid)
|
||||
|
||||
self.pool.remove(process)
|
||||
self.spawn_worker()
|
||||
if process.timer.value == 0:
|
||||
|
||||
@@ -2,7 +2,6 @@ from datetime import datetime, timedelta
|
||||
from keyword import iskeyword
|
||||
|
||||
# Django
|
||||
from django import get_version
|
||||
from django.core.exceptions import ValidationError
|
||||
from django.db import models
|
||||
from django.db.models import Q
|
||||
@@ -15,7 +14,6 @@ from django.utils.translation import gettext_lazy as _
|
||||
|
||||
# External
|
||||
from picklefield import PickledObjectField
|
||||
from picklefield.fields import dbsafe_decode
|
||||
|
||||
# Local
|
||||
from django_q.conf import croniter
|
||||
@@ -59,7 +57,7 @@ class Task(models.Model):
|
||||
.exclude(success=False)
|
||||
.values_list("result", flat=True)
|
||||
)
|
||||
return decode_results(values)
|
||||
return values
|
||||
|
||||
def group_result(self, failures=False):
|
||||
if self.group:
|
||||
@@ -298,12 +296,18 @@ class Schedule(models.Model):
|
||||
url = reverse("admin:django_q_success_change", args=(task.id,))
|
||||
else:
|
||||
url = reverse("admin:django_q_failure_change", args=(task.id,))
|
||||
return format_html(f'<a href="{url}">[{task.name}]</a>')
|
||||
return format_html('<a href="{}">[{}]</a>', url, task.name)
|
||||
return None
|
||||
|
||||
def __str__(self):
|
||||
return self.func
|
||||
|
||||
def save(self, *args, **kwargs):
|
||||
if self.pk is None and self.schedule_type == self.CRON:
|
||||
self.next_run = self.calculate_next_run()
|
||||
|
||||
return super().save(*args, **kwargs)
|
||||
|
||||
success.boolean = True
|
||||
success.short_description = _("success")
|
||||
last_run.allow_tags = True
|
||||
@@ -356,11 +360,3 @@ class OrmQ(models.Model):
|
||||
app_label = "django_q"
|
||||
verbose_name = _("Queued task")
|
||||
verbose_name_plural = _("Queued tasks")
|
||||
|
||||
|
||||
# Backwards compatibility for Django 1.7
|
||||
def decode_results(values):
|
||||
if get_version().split(".")[1] == "7":
|
||||
# decode values in 1.7
|
||||
return [dbsafe_decode(v) for v in values]
|
||||
return values
|
||||
|
||||
@@ -107,13 +107,16 @@ def save_task(task, broker: Broker):
|
||||
value = get_func_repr(value)
|
||||
filters[Conf.SAVE_LIMIT_PER] = value
|
||||
|
||||
with db.transaction.atomic(using=db.router.db_for_write(Success)):
|
||||
list(Success.objects.filter(**filters).select_for_update())
|
||||
if (
|
||||
task["success"]
|
||||
and 0 < Conf.SAVE_LIMIT <= Success.objects.filter(**filters).count()
|
||||
):
|
||||
Success.objects.filter(**filters).last().delete()
|
||||
# check if we should clean the success tasks
|
||||
if Conf.SAVE_LIMIT > 0:
|
||||
with db.transaction.atomic(using=db.router.db_for_write(Success)):
|
||||
success_tasks_qs = Success.objects.filter(**filters)
|
||||
success_tasks_pks = [
|
||||
success_task.pk
|
||||
for success_task in success_tasks_qs.select_for_update()
|
||||
]
|
||||
if task["success"] and len(success_tasks_pks) >= Conf.SAVE_LIMIT:
|
||||
success_tasks_qs.last().delete()
|
||||
|
||||
# check if this task has previous results
|
||||
try:
|
||||
@@ -145,7 +148,11 @@ def save_task(task, broker: Broker):
|
||||
attempt_count=1,
|
||||
)
|
||||
|
||||
if Conf.MAX_ATTEMPTS > 0 and task_obj.attempt_count >= Conf.MAX_ATTEMPTS:
|
||||
if (
|
||||
Conf.MAX_ATTEMPTS > 0
|
||||
and task_obj.attempt_count >= Conf.MAX_ATTEMPTS
|
||||
and task.get("ack_id")
|
||||
):
|
||||
broker.acknowledge(task["ack_id"])
|
||||
|
||||
except Exception:
|
||||
|
||||
@@ -65,7 +65,7 @@ def scheduler(broker: Broker = None):
|
||||
if s.args:
|
||||
args = ast.literal_eval(s.args)
|
||||
# single value won't eval to tuple, so:
|
||||
if type(args) != tuple:
|
||||
if type(args) is not tuple:
|
||||
args = (args,)
|
||||
q_options = kwargs.get("q_options", {})
|
||||
if s.intended_date_kwarg:
|
||||
|
||||
@@ -51,7 +51,7 @@ def test_redis(monkeypatch):
|
||||
|
||||
|
||||
def test_custom(monkeypatch):
|
||||
monkeypatch.setattr(Conf, "BROKER_CLASS", "brokers.redis_broker.Redis")
|
||||
monkeypatch.setattr(Conf, "BROKER_CLASS", "django_q.brokers.redis_broker.Redis")
|
||||
broker = get_broker()
|
||||
assert broker.ping() is True
|
||||
assert broker.info() is not None
|
||||
@@ -124,7 +124,7 @@ def test_ironmq(monkeypatch):
|
||||
@pytest.mark.skipif(
|
||||
not os.getenv("AWS_ACCESS_KEY_ID"), reason="requires AWS credentials"
|
||||
)
|
||||
def canceled_sqs(monkeypatch):
|
||||
def test_sqs(monkeypatch):
|
||||
monkeypatch.setattr(
|
||||
Conf,
|
||||
"SQS",
|
||||
@@ -132,11 +132,13 @@ def canceled_sqs(monkeypatch):
|
||||
"aws_region": os.getenv("AWS_REGION"),
|
||||
"aws_access_key_id": os.getenv("AWS_ACCESS_KEY_ID"),
|
||||
"aws_secret_access_key": os.getenv("AWS_SECRET_ACCESS_KEY"),
|
||||
"receive_message_wait_time_seconds": 20,
|
||||
"receive_message_wait_time_seconds": 5,
|
||||
},
|
||||
)
|
||||
# check broker
|
||||
broker = get_broker(list_key=uuid()[0])
|
||||
broker = get_broker(list_key="testing")
|
||||
assert "receive_message_wait_time_seconds" in Conf.SQS
|
||||
assert "aws_region" in Conf.SQS
|
||||
assert broker.ping() is True
|
||||
assert broker.info() is not None
|
||||
assert broker.queue_size() == 0
|
||||
@@ -173,7 +175,7 @@ def canceled_sqs(monkeypatch):
|
||||
broker.enqueue("test")
|
||||
while task is None:
|
||||
task = broker.dequeue()[0]
|
||||
broker.fail(task[0])
|
||||
broker.fail(task[0][0])
|
||||
# bulk test
|
||||
for _ in range(10):
|
||||
broker.enqueue("test")
|
||||
|
||||
@@ -399,7 +399,7 @@ def test_timeout_task_finishes(broker, cluster_config_timeout, async_task_kwargs
|
||||
|
||||
|
||||
@pytest.mark.django_db
|
||||
def test_recycle(broker, monkeypatch):
|
||||
def test_recycle(broker, monkeypatch, django_assert_num_queries):
|
||||
# set up the Sentinel
|
||||
broker.list_key = "test_recycle_test:q"
|
||||
async_task("django_q.tests.tasks.multiply", 2, 2, broker=broker)
|
||||
@@ -432,7 +432,8 @@ def test_recycle(broker, monkeypatch):
|
||||
monkeypatch.setattr(Conf, "SAVE_LIMIT", 1)
|
||||
result_queue.put("STOP")
|
||||
# run monitor
|
||||
monitor(result_queue)
|
||||
with django_assert_num_queries(12):
|
||||
monitor(result_queue)
|
||||
assert Success.objects.count() == Conf.SAVE_LIMIT
|
||||
broker.delete_queue()
|
||||
|
||||
@@ -487,7 +488,7 @@ def test_max_rss(broker, monkeypatch):
|
||||
stop_event = Event()
|
||||
cluster_id = uuidlib.uuid4()
|
||||
# override settings
|
||||
monkeypatch.setattr(Conf, "MAX_RSS", 40000)
|
||||
monkeypatch.setattr(Conf, "MAX_RSS", 20000)
|
||||
monkeypatch.setattr(Conf, "WORKERS", 1)
|
||||
# set a timer to stop the Sentinel
|
||||
threading.Timer(3, stop_event.set).start()
|
||||
|
||||
@@ -473,6 +473,47 @@ def test_scheduler_atomic_must_specify_the_database_based_on_router_redirection(
|
||||
mocked_db.atomic.assert_called_with(using="default")
|
||||
|
||||
|
||||
@pytest.mark.django_db
|
||||
def test_schedule_save_sets_next_run_for_cron():
|
||||
"""Ensure Schedule.save() sets next_run correctly for CRON schedules."""
|
||||
cron_expression = "0 12 * * *" # Executes dialy at 12pm
|
||||
schedule = Schedule(
|
||||
func="math.sqrt",
|
||||
schedule_type=Schedule.CRON,
|
||||
cron=cron_expression,
|
||||
)
|
||||
|
||||
assert schedule.next_run is not None
|
||||
assert schedule.pk is None
|
||||
|
||||
initial_next_run = schedule.next_run
|
||||
schedule.save()
|
||||
|
||||
# After save, next_run must be recalculated based on the CRON expression
|
||||
assert schedule.next_run > initial_next_run
|
||||
|
||||
|
||||
@pytest.mark.django_db
|
||||
def test_schedule_save_direct_db(broker):
|
||||
"""Ensure Schedule.save() updates next_run correctly when created directly in DB."""
|
||||
cron_expression = "0 12 * * *" # Executes dialy at 12pm
|
||||
|
||||
# Creating schedule directly in database
|
||||
schedule = Schedule.objects.create(
|
||||
func="math.sqrt",
|
||||
schedule_type=Schedule.CRON,
|
||||
cron=cron_expression,
|
||||
)
|
||||
|
||||
# Next run must be defined after execution time
|
||||
assert schedule.next_run is not None
|
||||
assert schedule.next_run > timezone.now()
|
||||
|
||||
scheduler(broker)
|
||||
|
||||
assert broker.queue_size() == 0
|
||||
|
||||
|
||||
def test_localtime():
|
||||
assert not is_naive(localtime())
|
||||
|
||||
|
||||
@@ -22,11 +22,13 @@ class TimeoutHandler:
|
||||
return
|
||||
try:
|
||||
signal.signal(signal.SIGALRM, self.raise_timeout_exception)
|
||||
except ValueError: # ValueError is raised for Windows users
|
||||
signal.alarm(self._timeout)
|
||||
except (
|
||||
ValueError,
|
||||
AttributeError,
|
||||
): # AttributeError or ValueError might be raised for Windows users
|
||||
logger.debug(_("SIGALARM is not available on your platform"))
|
||||
|
||||
signal.alarm(self._timeout)
|
||||
|
||||
def __exit__(self, exc_type, exc_value, traceback):
|
||||
if self._timeout == -1:
|
||||
return
|
||||
@@ -34,5 +36,8 @@ class TimeoutHandler:
|
||||
try:
|
||||
signal.alarm(0)
|
||||
signal.signal(signal.SIGALRM, signal.SIG_DFL)
|
||||
except ValueError: # ValueError is raised for Windows users
|
||||
except (
|
||||
ValueError,
|
||||
AttributeError,
|
||||
): # AttributeError or ValueError might be raised for Windows users
|
||||
logger.debug(_("SIGALARM is not available on your platform"))
|
||||
|
||||
@@ -73,9 +73,9 @@ author = "Ilan Steemers, Stan Triepels"
|
||||
# built documents.
|
||||
#
|
||||
# The short X.Y version.
|
||||
version = "1.7"
|
||||
version = "1.8"
|
||||
# The full version, including alpha/beta/rc tags.
|
||||
release = "1.7.0"
|
||||
release = "1.8.0"
|
||||
|
||||
# The language for content autogenerated by Sphinx. Refer to documentation
|
||||
# for a list of supported languages.
|
||||
|
||||
@@ -4,7 +4,7 @@ build-backend = "poetry.core.masonry.api"
|
||||
|
||||
[tool.poetry]
|
||||
name = "django-q2"
|
||||
version = "1.7.0"
|
||||
version = "1.8.0"
|
||||
packages = [
|
||||
{ include = "django_q" },
|
||||
]
|
||||
@@ -30,16 +30,19 @@ classifiers = [
|
||||
'Operating System :: MacOS',
|
||||
'Programming Language :: Python',
|
||||
'Programming Language :: Python :: 3',
|
||||
'Programming Language :: Python :: 3.8',
|
||||
'Programming Language :: Python :: 3.9',
|
||||
'Programming Language :: Python :: 3.10',
|
||||
'Programming Language :: Python :: 3.11',
|
||||
'Programming Language :: Python :: 3.12',
|
||||
'Programming Language :: Python :: 3.13',
|
||||
'Topic :: Internet :: WWW/HTTP',
|
||||
'Topic :: System :: Distributed Computing',
|
||||
'Topic :: Software Development :: Libraries :: Python Modules',
|
||||
]
|
||||
include = ['CHANGELOG.md']
|
||||
|
||||
include = [
|
||||
{ path = "CHANGELOG.md", format = "sdist" },
|
||||
]
|
||||
|
||||
[tool.poetry.plugins] # Optional super table
|
||||
|
||||
@@ -49,7 +52,7 @@ include = ['CHANGELOG.md']
|
||||
|
||||
|
||||
[tool.poetry.dependencies]
|
||||
python = ">=3.8,<4"
|
||||
python = ">=3.9,<4"
|
||||
django = ">=4.2, <6"
|
||||
django-picklefield = "^3.1"
|
||||
|
||||
|
||||
@@ -13,15 +13,41 @@ services:
|
||||
networks:
|
||||
- main
|
||||
|
||||
aws:
|
||||
container_name: aws
|
||||
image: localstack/localstack:3.4.0
|
||||
ports:
|
||||
- "127.0.0.1:4566:4566" # LocalStack Gateway
|
||||
- "127.0.0.1:4510-4559:4510-4559" # External services port range
|
||||
environment:
|
||||
AWS_DEFAULT_REGION: ${AWS_DEFAULT_REGION:-us-west-2}
|
||||
DEFAULT_REGION: ${AWS_DEFAULT_REGION:-us-west-2}
|
||||
SQS_ENDPOINT_STRATEGY: path
|
||||
SERVICES: sqs
|
||||
LOCALSTACK_HOST: aws
|
||||
DEBUG: 1
|
||||
LS_LOG: trace
|
||||
volumes:
|
||||
- ./containers/localstack:/etc/localstack/init/ready.d
|
||||
networks:
|
||||
- main
|
||||
|
||||
django-q2:
|
||||
build:
|
||||
dockerfile: ./Dockerfile.dev
|
||||
context: .
|
||||
environment:
|
||||
AWS_ENDPOINT_URL: http://aws:4566
|
||||
AWS_REGION: ${AWS_REGION:-us-west-2}
|
||||
AWS_ACCESS_KEY_ID: ${AWS_ACCESS_KEY_ID:-test}
|
||||
AWS_SECRET_ACCESS_KEY: ${AWS_SECRET_ACCESS_KEY:-test}
|
||||
AWS_DEFAULT_REGION: ${AWS_DEFAULT_REGION:-us-west-2}
|
||||
volumes:
|
||||
- .:/app
|
||||
depends_on:
|
||||
- redis
|
||||
- mongo
|
||||
- aws
|
||||
networks:
|
||||
- main
|
||||
|
||||
|
||||
Reference in New Issue
Block a user