Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions .pre-commit-config.yaml
Original file line number Diff line number Diff line change
@@ -1,6 +1,13 @@
exclude: |
(?x)
# NOT INSTALLABLE ADDONS
^base_import_async/|
^queue_job_batch/|
^queue_job_cron/|
^queue_job_cron_jobrunner/|
^queue_job_subscribe/|
^test_queue_job/|
^test_queue_job_batch/|
# END NOT INSTALLABLE ADDONS
# Files and folders generated by bots, to avoid loops
^setup/|/static/description/index\.html$|
Expand Down
10 changes: 5 additions & 5 deletions queue_job/README.rst
Original file line number Diff line number Diff line change
Expand Up @@ -21,13 +21,13 @@ Job Queue
:target: http://www.gnu.org/licenses/lgpl-3.0-standalone.html
:alt: License: LGPL-3
.. |badge3| image:: https://img.shields.io/badge/github-OCA%2Fqueue-lightgray.png?logo=github
:target: https://github.com/OCA/queue/tree/19.0/queue_job
:target: https://github.com/OCA/queue/tree/20.0/queue_job
:alt: OCA/queue
.. |badge4| image:: https://img.shields.io/badge/weblate-Translate%20me-F47D42.png
:target: https://translation.odoo-community.org/projects/queue-19-0/queue-19-0-queue_job
:target: https://translation.odoo-community.org/projects/queue-20-0/queue-20-0-queue_job
:alt: Translate me on Weblate
.. |badge5| image:: https://img.shields.io/badge/runboat-Try%20me-875A7B.png
:target: https://runboat.odoo-community.org/builds?repo=OCA/queue&target_branch=19.0
:target: https://runboat.odoo-community.org/builds?repo=OCA/queue&target_branch=20.0
:alt: Try me on Runboat

|badge1| |badge2| |badge3| |badge4| |badge5|
Expand Down Expand Up @@ -685,7 +685,7 @@ Bug Tracker
Bugs are tracked on `GitHub Issues <https://github.com/OCA/queue/issues>`_.
In case of trouble, please check there if your issue has already been reported.
If you spotted it first, help us to smash it by providing a detailed and welcomed
`feedback <https://github.com/OCA/queue/issues/new?body=module:%20queue_job%0Aversion:%2019.0%0A%0A**Steps%20to%20reproduce**%0A-%20...%0A%0A**Current%20behavior**%0A%0A**Expected%20behavior**>`_.
`feedback <https://github.com/OCA/queue/issues/new?body=module:%20queue_job%0Aversion:%2020.0%0A%0A**Steps%20to%20reproduce**%0A-%20...%0A%0A**Current%20behavior**%0A%0A**Expected%20behavior**>`_.

Do not contact contributors directly about support or help with technical issues.

Expand Down Expand Up @@ -747,6 +747,6 @@ Current `maintainers <https://odoo-community.org/page/maintainer-role>`__:

|maintainer-guewen| |maintainer-sbidoul|

This module is part of the `OCA/queue <https://github.com/OCA/queue/tree/19.0/queue_job>`_ project on GitHub.
This module is part of the `OCA/queue <https://github.com/OCA/queue/tree/20.0/queue_job>`_ project on GitHub.

You are welcome to contribute. To learn how please visit https://odoo-community.org/page/Contribute.
6 changes: 3 additions & 3 deletions queue_job/__manifest__.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

{
"name": "Job Queue",
"version": "19.0.2.1.3",
"version": "20.0.1.0.0",
"author": "Camptocamp,ACSONE SA/NV,Odoo Community Association (OCA)",
"website": "https://github.com/OCA/queue",
"license": "LGPL-3",
Expand All @@ -11,7 +11,7 @@
"external_dependencies": {"python": ["requests", "openupgradelib"]},
"data": [
"security/security.xml",
"security/ir.model.access.csv",
"security/ir.access.csv",
"views/queue_job_views.xml",
"views/queue_job_channel_views.xml",
"views/queue_job_function_views.xml",
Expand All @@ -27,7 +27,7 @@
"queue_job/static/src/views/**/*",
],
},
"installable": False,
"installable": True,
"development_status": "Mature",
"maintainers": ["guewen", "sbidoul"],
"post_init_hook": "post_init_hook",
Expand Down
15 changes: 7 additions & 8 deletions queue_job/controllers/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@
from werkzeug.exceptions import BadRequest, Forbidden

from odoo import SUPERUSER_ID, api, http
from odoo.service.model import PG_CONCURRENCY_ERRORS_TO_RETRY
from odoo.sql_db import PG_CONCURRENCY_EXCEPTIONS_TO_RETRY
from odoo.tools import config

from ..delay import chain, group
Expand Down Expand Up @@ -123,7 +123,7 @@ def _enqueue_dependent_jobs(cls, env, job):
except OperationalError as err:
# Automatically retry the typical transaction serialization
# errors
if err.pgcode not in PG_CONCURRENCY_ERRORS_TO_RETRY:
if not isinstance(err, PG_CONCURRENCY_EXCEPTIONS_TO_RETRY):
raise
if tries >= DEPENDS_MAX_TRIES_ON_CONCURRENCY_FAILURE:
_logger.error(
Expand All @@ -148,7 +148,7 @@ def _enqueue_dependent_jobs(cls, env, job):
@classmethod
def _runjob(cls, env: api.Environment, job: Job) -> None:
def retry_postpone(job, message, seconds=None):
job.env.clear()
job.env.transaction.clear()
with job.in_temporary_env():
job.postpone(result=message, seconds=seconds)
job.set_pending(reset_retry=False)
Expand All @@ -160,7 +160,7 @@ def retry_postpone(job, message, seconds=None):
except OperationalError as err:
# Automatically retry the typical transaction serialization
# errors
if err.pgcode not in PG_CONCURRENCY_ERRORS_TO_RETRY:
if not isinstance(err, PG_CONCURRENCY_EXCEPTIONS_TO_RETRY):
raise

_logger.debug("%s OperationalError, postponed", job)
Expand All @@ -181,7 +181,7 @@ def retry_postpone(job, message, seconds=None):
traceback.print_exc(file=buff)
traceback_txt = buff.getvalue()
_logger.error(traceback_txt)
job.env.clear()
job.env.transaction.clear()
with job.in_temporary_env():
vals = cls._get_failure_values(job, traceback_txt, orig_exception)
job.set_failed(**vals)
Expand Down Expand Up @@ -388,6 +388,5 @@ def _create_graph_test_jobs(

root_delayable.delay()

return (
f"graph uuid: {list(root_delayable._head())[0]._generated_job.graph_uuid}"
)
head = next(iter(root_delayable._head()))
return f"graph uuid: {head._generated_job.graph_uuid}"
6 changes: 3 additions & 3 deletions queue_job/data/queue_data.xml
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
<?xml version="1.0" encoding="utf-8" ?>
<?xml version="1.0" encoding="UTF-8" ?>
<odoo>
<data noupdate="1">
<!-- Queue-job-related subtypes for messaging / Chatter -->
Expand All @@ -12,14 +12,14 @@
<field ref="model_queue_job" name="model_id" />
<field eval="True" name="active" />
<field name="user_id" ref="base.user_root" />
<field name="interval_number">1</field>
<field name="interval_number" eval="1" />
<field name="interval_type">days</field>
<field name="state">code</field>
<field name="code">model.autovacuum()</field>
</record>
</data>
<data noupdate="0">
<record model="queue.job.channel" id="channel_root">
<record id="channel_root" model="queue.job.channel">
<field name="name">root</field>
</record>
</data>
Expand Down
1 change: 1 addition & 0 deletions queue_job/data/queue_job_function_data.xml
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
<?xml version="1.0" encoding="UTF-8" ?>
<odoo noupdate="1">
<record id="job_function_queue_job__test_job" model="queue.job.function">
<field name="model_id" ref="queue_job.model_queue_job" />
Expand Down
6 changes: 3 additions & 3 deletions queue_job/delay.py
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,7 @@ class Graph:
instances, although ultimately it is used for this purpose.
"""

__slots__ = "_graph"
__slots__ = ("_graph",)

def __init__(self, graph=None):
if graph:
Expand Down Expand Up @@ -314,7 +314,7 @@ class DelayableChain:
delayable/chain/group object of the graph.
"""

__slots__ = ("_graph", "__head", "__tail")
__slots__ = ("__head", "__tail", "_graph")

def __init__(self, *delayables):
self._graph = DelayableGraph()
Expand Down Expand Up @@ -371,7 +371,7 @@ class DelayableGroup:
delayable/chain/group object of the graph.
"""

__slots__ = ("_graph", "_delayables")
__slots__ = ("_delayables", "_graph")

def __init__(self, *delayables):
self._graph = DelayableGraph()
Expand Down
6 changes: 3 additions & 3 deletions queue_job/fields.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ class JobSerialized(fields.Json):
_base_type = None

# these are the default values when we convert an empty value
_default_json_mapping = {
_default_json_mapping = { # noqa: RUF012
dict: "{}",
list: "[]",
tuple: "[]",
Expand All @@ -42,8 +42,8 @@ class JobSerialized(fields.Json):
def __init__(self, string=SENTINEL, base_type=SENTINEL, **kwargs):
super().__init__(string=string, _base_type=base_type, **kwargs)

def _setup_attrs(self, model, name): # pylint: disable=missing-return
super()._setup_attrs(model, name)
def _setup_attrs__(self, model_class, name): # pylint: disable=missing-return
super()._setup_attrs__(model_class, name)
if self._base_type not in self._default_json_mapping:
msg = f"{self._base_type} is not a supported base type"
raise ValueError(msg)
Expand Down
23 changes: 11 additions & 12 deletions queue_job/job.py
Original file line number Diff line number Diff line change
Expand Up @@ -461,7 +461,7 @@ def __init__(
if self.priority is None:
self.priority = DEFAULT_PRIORITY

self.date_created = datetime.now()
self.date_created = datetime.now() # noqa: DTZ005
self._description = description

if isinstance(identity_key, str):
Expand Down Expand Up @@ -524,7 +524,7 @@ def perform(self):
elif not self.max_retries: # infinite retries
raise
elif self.retry >= self.max_retries:
type_, value, traceback = sys.exc_info()
type_, value, _traceback = sys.exc_info()
# change the exception type but keep the original
# traceback and message:
# http://blog.ianbicking.org/2007/09/12/re-raising-exceptions/
Expand Down Expand Up @@ -745,9 +745,8 @@ def job_function_name(self):

@property
def identity_key(self):
if self._identity_key is None:
if self._identity_key_func:
self._identity_key = self._identity_key_func(self)
if self._identity_key is None and self._identity_key_func:
self._identity_key = self._identity_key_func(self)
return self._identity_key

@identity_key.setter
Expand Down Expand Up @@ -808,9 +807,9 @@ def eta(self, value):
if not value:
self._eta = None
elif isinstance(value, timedelta):
self._eta = datetime.now() + value
self._eta = datetime.now() + value # noqa: DTZ005
elif isinstance(value, int):
self._eta = datetime.now() + timedelta(seconds=value)
self._eta = datetime.now() + timedelta(seconds=value) # noqa: DTZ005
else:
self._eta = value

Expand Down Expand Up @@ -845,27 +844,27 @@ def set_pending(self, result=None, reset_retry=True):

def set_enqueued(self):
self.state = ENQUEUED
self.date_enqueued = datetime.now()
self.date_enqueued = datetime.now() # noqa: DTZ005
self.date_started = None
self.worker_pid = None

def set_started(self):
self.state = STARTED
self.date_started = datetime.now()
self.date_started = datetime.now() # noqa: DTZ005
self.worker_pid = os.getpid()
self.add_lock_record()

def set_done(self, result=None):
self.state = DONE
self.exc_name = None
self.exc_info = None
self.date_done = datetime.now()
self.date_done = datetime.now() # noqa: DTZ005
if result is not None:
self.result = result

def set_cancelled(self, result=None):
self.state = CANCELLED
self.date_cancelled = datetime.now()
self.date_cancelled = datetime.now() # noqa: DTZ005
if result is not None:
self.result = result

Expand Down Expand Up @@ -925,7 +924,7 @@ def related_action(self):
if not funcname:
funcname = record._default_related_action
if not isinstance(funcname, str):
raise ValueError(
raise ValueError( # noqa: TRY004
"related_action must be the name of the method on queue.job as string"
)
action = getattr(record, funcname)
Expand Down
2 changes: 1 addition & 1 deletion queue_job/jobrunner/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@
queue_job_config = dict(cp["queue_job"])


from .runner import QueueJobRunner, _channels # noqa: E402
from .runner import QueueJobRunner, _channels

_logger = logging.getLogger(__name__)

Expand Down
2 changes: 1 addition & 1 deletion queue_job/jobrunner/channels.py
Original file line number Diff line number Diff line change
Expand Up @@ -179,7 +179,7 @@ class ChannelJob:

"""

__slots__ = ("db_name", "channel", "uuid", "_sorting_key", "__weakref__")
__slots__ = ("__weakref__", "_sorting_key", "channel", "db_name", "uuid")

def __init__(self, db_name, channel, uuid, seq, date_created, priority, eta):
self.db_name = db_name
Expand Down
28 changes: 14 additions & 14 deletions queue_job/jobrunner/runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
from psycopg2.extensions import ISOLATION_LEVEL_AUTOCOMMIT

import odoo
from odoo.modules.db import list_dbs
from odoo.tools import config

from . import queue_job_config
Expand Down Expand Up @@ -72,7 +73,7 @@ def _odoo_now():


def _connection_info_for(db_name):
db_or_uri, connection_info = odoo.sql_db.connection_info_for(db_name)
_db_or_uri, connection_info = odoo.sql_db.connection_info_for(db_name)

for p in ("host", "port", "user", "password"):
cfg = os.environ.get(
Expand Down Expand Up @@ -139,7 +140,7 @@ def close(self):
# del
try:
self.conn.close()
except Exception:
except Exception: # noqa: BLE001, S110
pass
self.conn = None

Expand Down Expand Up @@ -382,7 +383,7 @@ def get_db_names(self):
db_names = config["db_name"]
if db_names:
return db_names
return odoo.service.db.list_dbs(True)
return list_dbs(force=True)

def close_databases(self, remove_jobs=True):
for db_name, db in self.db_by_name.items():
Expand Down Expand Up @@ -474,17 +475,16 @@ def wait_notification(self):
# if timeout remains a large negative number, it is most
# probably a bug
_logger.debug("select() timeout: %.2f sec", timeout)
if timeout > 0:
if conns and not self._stop:
with select() as sel:
for conn in conns:
sel.register(conn, selectors.EVENT_READ)
events = sel.select(timeout=timeout)
for key, _mask in events:
if key.fileobj == self._stop_pipe[0]:
# stop-pipe is not a conn so doesn't need poll()
continue
key.fileobj.poll()
if timeout > 0 and conns and not self._stop:
with select() as sel:
for conn in conns:
sel.register(conn, selectors.EVENT_READ)
events = sel.select(timeout=timeout)
for key, _mask in events:
if key.fileobj == self._stop_pipe[0]:
# stop-pipe is not a conn so doesn't need poll()
continue
key.fileobj.poll()

def stop(self):
_logger.info("graceful stop requested")
Expand Down
Loading
Loading