mirror of
https://github.com/django-q2/django-q2.git
synced 2026-09-28 16:28:12 +08:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ccd760a69f | ||
|
|
fb23b6b825 | ||
|
|
35c911080b | ||
|
|
81058fc1cf | ||
|
|
505eb8d537 | ||
|
|
dcfe5b6650 | ||
|
|
c3b49cca25 | ||
|
|
c1e00b09cb | ||
|
|
8e403d0136 | ||
|
|
337a5782d9 | ||
|
|
2c456cd0aa | ||
|
|
8f54e5d5ae | ||
|
|
689c881178 |
@@ -15,6 +15,12 @@ jobs:
|
|||||||
run: |
|
run: |
|
||||||
pipx install ruff==0.4.10
|
pipx install ruff==0.4.10
|
||||||
ruff format . --check && ruff check .
|
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:
|
test:
|
||||||
runs-on: ubuntu-latest
|
runs-on: ubuntu-latest
|
||||||
@@ -93,11 +99,7 @@ jobs:
|
|||||||
- name: Upload to coveralls
|
- name: Upload to coveralls
|
||||||
run: |
|
run: |
|
||||||
python -m pip install --upgrade pip
|
python -m pip install --upgrade pip
|
||||||
python -m pip install coveralls flake8 black
|
python -m pip install coveralls
|
||||||
coveralls --service=github --finish
|
coveralls --service=github --finish
|
||||||
env:
|
env:
|
||||||
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||||
- name: Check flake8/black
|
|
||||||
run: |
|
|
||||||
flake8 .
|
|
||||||
black --check .
|
|
||||||
|
|||||||
+25
-1
@@ -1,6 +1,30 @@
|
|||||||
# Changelog
|
# Changelog
|
||||||
|
|
||||||
## [v1.7.0](https://github.com/django-q2/django-q2/tree/v1.7.0) (2024-08-24)
|
## [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:**
|
**Merged pull requests:**
|
||||||
|
|
||||||
|
|||||||
@@ -7,6 +7,10 @@ ENV PYTHONUNBUFFERED 1
|
|||||||
# Sets the default shell to bash
|
# Sets the default shell to bash
|
||||||
ENV SHELL /bin/bash
|
ENV SHELL /bin/bash
|
||||||
|
|
||||||
|
RUN set -ex \
|
||||||
|
&& apt update \
|
||||||
|
&& apt-get install gcc python3-dev --yes
|
||||||
|
|
||||||
# Upgrades pip
|
# Upgrades pip
|
||||||
RUN pip install -U pip setuptools
|
RUN pip install -U pip setuptools
|
||||||
|
|
||||||
|
|||||||
+2
-7
@@ -1,7 +1,7 @@
|
|||||||
A multiprocessing distributed task queue for Django
|
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
|
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
|
||||||
|
|
||||||
@@ -198,7 +198,7 @@ Admin page or directly from your code:
|
|||||||
For more info check the `Schedules <https://django-q2.readthedocs.org/en/latest/schedules.html>`__ documentation.
|
For more info check the `Schedules <https://django-q2.readthedocs.org/en/latest/schedules.html>`__ documentation.
|
||||||
|
|
||||||
Development
|
Development
|
||||||
~~~~~~~
|
~~~~~~~~~~~
|
||||||
|
|
||||||
There is an example project that you can use to develop with. Docker (compose) is being used to set everything up.
|
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.
|
Please note that you will have to restart the django-q container when changes have been made to tasks or django-q.
|
||||||
@@ -214,7 +214,6 @@ Create a superuser with:
|
|||||||
|
|
||||||
make createsuperuser
|
make createsuperuser
|
||||||
|
|
||||||
|
|
||||||
Testing
|
Testing
|
||||||
~~~~~~~
|
~~~~~~~
|
||||||
|
|
||||||
@@ -246,9 +245,5 @@ Acknowledgements
|
|||||||
:target: https://github.com/GDay/django-q2/actions?query=workflow%3Atests
|
: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
|
.. |image1| image:: https://coveralls.io/repos/github/GDay/django-q2/badge.svg?branch=master
|
||||||
:target: https://coveralls.io/github/GDay/django-q2?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
|
.. |downloads| image:: https://img.shields.io/pypi/dm/django-q2
|
||||||
:target: https://img.shields.io/pypi/dm/django-q2
|
:target: https://img.shields.io/pypi/dm/django-q2
|
||||||
|
|||||||
Executable
+19
@@ -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,6 +1,6 @@
|
|||||||
import django
|
import django
|
||||||
|
|
||||||
VERSION = (1, 7, 0)
|
VERSION = (1, 7, 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"
|
||||||
|
|||||||
@@ -1,3 +1,5 @@
|
|||||||
|
import copy
|
||||||
|
|
||||||
from boto3 import Session
|
from boto3 import Session
|
||||||
from botocore.client import ClientError
|
from botocore.client import ClientError
|
||||||
|
|
||||||
@@ -78,15 +80,15 @@ class Sqs(Broker):
|
|||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def get_connection(list_key: str = None) -> Session:
|
def get_connection(list_key: str = None) -> Session:
|
||||||
config = Conf.SQS
|
config_cloned = copy.deepcopy(Conf.SQS)
|
||||||
if "aws_region" in config:
|
if "aws_region" in config_cloned:
|
||||||
config["region_name"] = config["aws_region"]
|
config_cloned["region_name"] = config_cloned["aws_region"]
|
||||||
del config["aws_region"]
|
del config_cloned["aws_region"]
|
||||||
|
|
||||||
if "receive_message_wait_time_seconds" in config:
|
if "receive_message_wait_time_seconds" in config_cloned:
|
||||||
del config["receive_message_wait_time_seconds"]
|
del config_cloned["receive_message_wait_time_seconds"]
|
||||||
|
|
||||||
return Session(**config)
|
return Session(**config_cloned)
|
||||||
|
|
||||||
def get_queue(self):
|
def get_queue(self):
|
||||||
self.sqs = self.connection.resource("sqs")
|
self.sqs = self.connection.resource("sqs")
|
||||||
|
|||||||
+8
-1
@@ -1,4 +1,5 @@
|
|||||||
# Standard
|
# Standard
|
||||||
|
import os
|
||||||
import signal
|
import signal
|
||||||
import socket
|
import socket
|
||||||
import uuid
|
import uuid
|
||||||
@@ -230,8 +231,14 @@ class Sentinel:
|
|||||||
% {"name": process.name}
|
% {"name": process.name}
|
||||||
)
|
)
|
||||||
else:
|
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)
|
prometheus_multiprocess.mark_process_dead(process.pid)
|
||||||
|
|
||||||
self.pool.remove(process)
|
self.pool.remove(process)
|
||||||
self.spawn_worker()
|
self.spawn_worker()
|
||||||
if process.timer.value == 0:
|
if process.timer.value == 0:
|
||||||
|
|||||||
+5
-1
@@ -145,7 +145,11 @@ def save_task(task, broker: Broker):
|
|||||||
attempt_count=1,
|
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"])
|
broker.acknowledge(task["ack_id"])
|
||||||
|
|
||||||
except Exception:
|
except Exception:
|
||||||
|
|||||||
@@ -65,7 +65,7 @@ def scheduler(broker: Broker = None):
|
|||||||
if s.args:
|
if s.args:
|
||||||
args = ast.literal_eval(s.args)
|
args = ast.literal_eval(s.args)
|
||||||
# single value won't eval to tuple, so:
|
# single value won't eval to tuple, so:
|
||||||
if type(args) != tuple:
|
if type(args) is not tuple:
|
||||||
args = (args,)
|
args = (args,)
|
||||||
q_options = kwargs.get("q_options", {})
|
q_options = kwargs.get("q_options", {})
|
||||||
if s.intended_date_kwarg:
|
if s.intended_date_kwarg:
|
||||||
|
|||||||
@@ -51,7 +51,7 @@ def test_redis(monkeypatch):
|
|||||||
|
|
||||||
|
|
||||||
def test_custom(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()
|
broker = get_broker()
|
||||||
assert broker.ping() is True
|
assert broker.ping() is True
|
||||||
assert broker.info() is not None
|
assert broker.info() is not None
|
||||||
@@ -124,7 +124,7 @@ def test_ironmq(monkeypatch):
|
|||||||
@pytest.mark.skipif(
|
@pytest.mark.skipif(
|
||||||
not os.getenv("AWS_ACCESS_KEY_ID"), reason="requires AWS credentials"
|
not os.getenv("AWS_ACCESS_KEY_ID"), reason="requires AWS credentials"
|
||||||
)
|
)
|
||||||
def canceled_sqs(monkeypatch):
|
def test_sqs(monkeypatch):
|
||||||
monkeypatch.setattr(
|
monkeypatch.setattr(
|
||||||
Conf,
|
Conf,
|
||||||
"SQS",
|
"SQS",
|
||||||
@@ -132,11 +132,13 @@ def canceled_sqs(monkeypatch):
|
|||||||
"aws_region": os.getenv("AWS_REGION"),
|
"aws_region": os.getenv("AWS_REGION"),
|
||||||
"aws_access_key_id": os.getenv("AWS_ACCESS_KEY_ID"),
|
"aws_access_key_id": os.getenv("AWS_ACCESS_KEY_ID"),
|
||||||
"aws_secret_access_key": os.getenv("AWS_SECRET_ACCESS_KEY"),
|
"aws_secret_access_key": os.getenv("AWS_SECRET_ACCESS_KEY"),
|
||||||
"receive_message_wait_time_seconds": 20,
|
"receive_message_wait_time_seconds": 5,
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
# check broker
|
# 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.ping() is True
|
||||||
assert broker.info() is not None
|
assert broker.info() is not None
|
||||||
assert broker.queue_size() == 0
|
assert broker.queue_size() == 0
|
||||||
@@ -173,7 +175,7 @@ def canceled_sqs(monkeypatch):
|
|||||||
broker.enqueue("test")
|
broker.enqueue("test")
|
||||||
while task is None:
|
while task is None:
|
||||||
task = broker.dequeue()[0]
|
task = broker.dequeue()[0]
|
||||||
broker.fail(task[0])
|
broker.fail(task[0][0])
|
||||||
# bulk test
|
# bulk test
|
||||||
for _ in range(10):
|
for _ in range(10):
|
||||||
broker.enqueue("test")
|
broker.enqueue("test")
|
||||||
|
|||||||
@@ -487,7 +487,7 @@ def test_max_rss(broker, monkeypatch):
|
|||||||
stop_event = Event()
|
stop_event = Event()
|
||||||
cluster_id = uuidlib.uuid4()
|
cluster_id = uuidlib.uuid4()
|
||||||
# override settings
|
# override settings
|
||||||
monkeypatch.setattr(Conf, "MAX_RSS", 40000)
|
monkeypatch.setattr(Conf, "MAX_RSS", 20000)
|
||||||
monkeypatch.setattr(Conf, "WORKERS", 1)
|
monkeypatch.setattr(Conf, "WORKERS", 1)
|
||||||
# set a timer to stop the Sentinel
|
# set a timer to stop the Sentinel
|
||||||
threading.Timer(3, stop_event.set).start()
|
threading.Timer(3, stop_event.set).start()
|
||||||
|
|||||||
+9
-4
@@ -22,11 +22,13 @@ class TimeoutHandler:
|
|||||||
return
|
return
|
||||||
try:
|
try:
|
||||||
signal.signal(signal.SIGALRM, self.raise_timeout_exception)
|
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"))
|
logger.debug(_("SIGALARM is not available on your platform"))
|
||||||
|
|
||||||
signal.alarm(self._timeout)
|
|
||||||
|
|
||||||
def __exit__(self, exc_type, exc_value, traceback):
|
def __exit__(self, exc_type, exc_value, traceback):
|
||||||
if self._timeout == -1:
|
if self._timeout == -1:
|
||||||
return
|
return
|
||||||
@@ -34,5 +36,8 @@ class TimeoutHandler:
|
|||||||
try:
|
try:
|
||||||
signal.alarm(0)
|
signal.alarm(0)
|
||||||
signal.signal(signal.SIGALRM, signal.SIG_DFL)
|
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"))
|
logger.debug(_("SIGALARM is not available on your platform"))
|
||||||
|
|||||||
+1
-1
@@ -75,7 +75,7 @@ author = "Ilan Steemers, Stan Triepels"
|
|||||||
# The short X.Y version.
|
# The short X.Y version.
|
||||||
version = "1.7"
|
version = "1.7"
|
||||||
# The full version, including alpha/beta/rc tags.
|
# The full version, including alpha/beta/rc tags.
|
||||||
release = "1.7.0"
|
release = "1.7.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.
|
||||||
|
|||||||
+1
-1
@@ -4,7 +4,7 @@ build-backend = "poetry.core.masonry.api"
|
|||||||
|
|
||||||
[tool.poetry]
|
[tool.poetry]
|
||||||
name = "django-q2"
|
name = "django-q2"
|
||||||
version = "1.7.0"
|
version = "1.7.4"
|
||||||
packages = [
|
packages = [
|
||||||
{ include = "django_q" },
|
{ include = "django_q" },
|
||||||
]
|
]
|
||||||
|
|||||||
@@ -13,15 +13,41 @@ services:
|
|||||||
networks:
|
networks:
|
||||||
- main
|
- 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:
|
django-q2:
|
||||||
build:
|
build:
|
||||||
dockerfile: ./Dockerfile.dev
|
dockerfile: ./Dockerfile.dev
|
||||||
context: .
|
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:
|
volumes:
|
||||||
- .:/app
|
- .:/app
|
||||||
depends_on:
|
depends_on:
|
||||||
- redis
|
- redis
|
||||||
- mongo
|
- mongo
|
||||||
|
- aws
|
||||||
networks:
|
networks:
|
||||||
- main
|
- main
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user