diff --git a/.oca/oca-port/blacklist/queue_job.json b/.oca/oca-port/blacklist/queue_job.json new file mode 100644 index 0000000000..3fd748577f --- /dev/null +++ b/.oca/oca-port/blacklist/queue_job.json @@ -0,0 +1,12 @@ +{ + "pull_requests": { + "OCA/queue#715": "Already in 14.0", + "OCA/queue#726": "Not applicable", + "OCA/queue#743": "Already in 14.0", + "OCA/queue#782": "Already in 14.0", + "OCA/queue#775": "Already in 14.0", + "OCA/queue#804": "Not applicable", + "OCA/queue#854": "(auto) Nothing to port from PR #854", + "OCA/queue#938": "Not applicable" + } +} diff --git a/queue_job/__manifest__.py b/queue_job/__manifest__.py index e25e4742e9..698e2fc8d3 100644 --- a/queue_job/__manifest__.py +++ b/queue_job/__manifest__.py @@ -2,7 +2,7 @@ { "name": "Job Queue", - "version": "14.0.3.15.1", + "version": "14.0.3.16.0", "author": "Camptocamp,ACSONE SA/NV,Odoo Community Association (OCA)", "website": "https://github.com/OCA/queue", "license": "LGPL-3", diff --git a/queue_job/controllers/main.py b/queue_job/controllers/main.py index ddec4d95ca..c4f7028a74 100644 --- a/queue_job/controllers/main.py +++ b/queue_job/controllers/main.py @@ -6,14 +6,16 @@ import random import time import traceback +from contextlib import contextmanager from io import StringIO +from typing import Union from psycopg2 import OperationalError, errorcodes from werkzeug.exceptions import BadRequest, Forbidden -import odoo -from odoo import _, http, tools +from odoo import SUPERUSER_ID, _, api, http, tools from odoo.service.model import PG_CONCURRENCY_ERRORS_TO_RETRY +from odoo.tools import config from ..delay import chain, group from ..exception import FailedJobError, NothingToDoJob, RetryableJobError @@ -26,36 +28,106 @@ DEPENDS_MAX_TRIES_ON_CONCURRENCY_FAILURE = 5 +@contextmanager +def _prevent_commit(cr): + """Context manager to prevent commits on a cursor. + + Commiting while the job is not finished would release the job lock, causing + it to be started again by the dead jobs requeuer. + """ + + def forbidden_commit(*args, **kwargs): + raise RuntimeError( + "Commit is forbidden in queue jobs. " + 'You may want to enable the "Allow Commit" option on the Job ' + "Function. Alternatively, if the current job is a cron running as " + "queue job, you can modify it to run as a normal cron. More details on: " + "https://github.com/OCA/queue/wiki/Upgrade-warning:-commits-inside-jobs" + ) + + original_commit = cr.commit + cr.commit = forbidden_commit + try: + yield + finally: + cr.commit = original_commit + + class RunJobController(http.Controller): - def _try_perform_job(self, env, job): - """Try to perform the job.""" + @classmethod + def _acquire_job(cls, env: api.Environment, job_uuid: str) -> Union[Job, None]: + """Acquire a job for execution. + + - make sure it is in ENQUEUED state + - mark it as STARTED and commit the state change + - acquire the job lock + + If successful, return the Job instance, otherwise return None. This + function may fail to acquire the job is not in the expected state or is + already locked by another worker. + """ + env.cr.execute( + "SELECT uuid FROM queue_job WHERE uuid=%s AND state=%s " + "FOR NO KEY UPDATE SKIP LOCKED", + (job_uuid, ENQUEUED), + ) + if not env.cr.fetchone(): + _logger.warning( + "was requested to run job %s, but it does not exist, " + "or is not in state %s, or is being handled by another worker", + job_uuid, + ENQUEUED, + ) + return None + job = Job.load(env, job_uuid) + assert job and job.state == ENQUEUED job.set_started() job.store() env.cr.commit() - _logger.debug("%s started", job) + if not job.lock(): + _logger.warning( + "was requested to run job %s, but it could not be locked", + job_uuid, + ) + return None + return job - job.perform() - # Triggers any stored computed fields before calling 'set_done' - # so that will be part of the 'exec_time' - env["base"].flush() - job.set_done() - job.store() - env["base"].flush() - env.cr.commit() + @classmethod + def _try_perform_job(cls, env, job): + """Try to perform the job, mark it done and commit if successful.""" + _logger.debug("%s started", job) + # TODO refactor, the relation between env and job.env is not clear + assert env.cr is job.env.cr + with _prevent_commit(env.cr): + job.perform() + # Triggers any stored computed fields before calling 'set_done' + # so that will be part of the 'exec_time' + env["base"].flush() + job.set_done() + job.store() + env["base"].flush() + if not config["test_enable"]: + env.cr.commit() _logger.debug("%s done", job) - def _enqueue_dependent_jobs(self, env, job): + @classmethod + def _enqueue_dependent_jobs(cls, env, job): + if not job.should_check_dependents(): + return + + _logger.debug("%s enqueue depends started", job) tries = 0 while True: try: - job.enqueue_waiting() + with job.env.cr.savepoint(): + job.enqueue_waiting() except OperationalError as err: # Automatically retry the typical transaction serialization # errors if err.pgcode not in PG_CONCURRENCY_ERRORS_TO_RETRY: raise if tries >= DEPENDS_MAX_TRIES_ON_CONCURRENCY_FAILURE: - _logger.info( + _logger.error( "%s, maximum number of tries reached to update dependencies", errorcodes.lookup(err.pgcode), ) @@ -72,42 +144,20 @@ def _enqueue_dependent_jobs(self, env, job): time.sleep(wait_time) else: break + _logger.debug("%s enqueue depends done", job) - @http.route("/queue_job/runjob", type="http", auth="none", save_session=False) - def runjob(self, db, job_uuid, **kw): - http.request.session.db = db - env = http.request.env(user=odoo.SUPERUSER_ID) - + @classmethod + def _runjob(cls, env: api.Environment, job: Job) -> None: def retry_postpone(job, message, seconds=None): job.env.clear() - with odoo.api.Environment.manage(): - with odoo.registry(job.env.cr.dbname).cursor() as new_cr: - job.env = job.env(cr=new_cr) - job.postpone(result=message, seconds=seconds) - job.set_pending(reset_retry=False) - job.store() - new_cr.commit() - - # ensure the job to run is in the correct state and lock the record - env.cr.execute( - "SELECT state FROM queue_job WHERE uuid=%s AND state=%s FOR UPDATE", - (job_uuid, ENQUEUED), - ) - if not env.cr.fetchone(): - _logger.warning( - "was requested to run job %s, but it does not exist, " - "or is not in state %s", - job_uuid, - ENQUEUED, - ) - return "" - - job = Job.load(env, job_uuid) - assert job and job.state == ENQUEUED + with job.in_temporary_env(): + job.postpone(result=message, seconds=seconds) + job.set_pending(reset_retry=False) + job.store() try: try: - self._try_perform_job(env, job) + cls._try_perform_job(env, job) except OperationalError as err: # Automatically retry the typical transaction serialization # errors @@ -136,7 +186,7 @@ def retry_postpone(job, message, seconds=None): # traceback in the logs we should have the traceback when all # retries are exhausted env.cr.rollback() - return "" + return except (FailedJobError, Exception) as orig_exception: buff = StringIO() @@ -144,34 +194,48 @@ def retry_postpone(job, message, seconds=None): traceback_txt = buff.getvalue() _logger.error(traceback_txt) job.env.clear() - with odoo.api.Environment.manage(): - with odoo.registry(job.env.cr.dbname).cursor() as new_cr: - job.env = job.env(cr=new_cr) - vals = self._get_failure_values(job, traceback_txt, orig_exception) - job.set_failed(**vals) - job.store() - new_cr.commit() - buff.close() + with job.in_temporary_env(): + vals = cls._get_failure_values(job, traceback_txt, orig_exception) + job.set_failed(**vals) + job.store() + job.on_fail(vals) + buff.close() raise - _logger.debug("%s enqueue depends started", job) - self._enqueue_dependent_jobs(env, job) - _logger.debug("%s enqueue depends done", job) - - return "" + cls._enqueue_dependent_jobs(env, job) - def _get_failure_values(self, job, traceback_txt, orig_exception): + @classmethod + def _get_failure_values(cls, job, traceback_txt, orig_exception): """Collect relevant data from exception.""" exception_name = orig_exception.__class__.__name__ if hasattr(orig_exception, "__module__"): exception_name = orig_exception.__module__ + "." + exception_name - exc_message = getattr(orig_exception, "name", str(orig_exception)) + exc_message = ( + orig_exception.args[0] if orig_exception.args else str(orig_exception) + ) return { "exc_info": traceback_txt, "exc_name": exception_name, "exc_message": exc_message, } + @http.route( + "/queue_job/runjob", + type="http", + auth="none", + save_session=False, + readonly=False, + ) + def runjob(self, db, job_uuid, **kw): + http.request.session.db = db + env = http.request.env(user=SUPERUSER_ID) + job = self._acquire_job(env, job_uuid) + if not job: + return "" + self._runjob(env, job) + return "" + + # flake8: noqa: C901 @http.route("/queue_job/create_test_job", type="http", auth="user") def create_test_job( self, @@ -181,6 +245,8 @@ def create_test_job( description="Test job", size=1, failure_rate=0, + commit_within_job=False, + failure_retry_seconds=0, ): """Create test jobs @@ -222,6 +288,12 @@ def create_test_job( except ValueError: max_retries = None + if failure_retry_seconds is not None: + try: + failure_retry_seconds = int(failure_retry_seconds) + except ValueError: + failure_retry_seconds = 0 + if size == 1: return self._create_single_test_job( priority=priority, @@ -229,6 +301,8 @@ def create_test_job( channel=channel, description=description, failure_rate=failure_rate, + commit_within_job=commit_within_job, + failure_retry_seconds=failure_retry_seconds, ) if size > 1: @@ -239,6 +313,8 @@ def create_test_job( channel=channel, description=description, failure_rate=failure_rate, + commit_within_job=commit_within_job, + failure_retry_seconds=failure_retry_seconds, ) return "" @@ -250,6 +326,8 @@ def _create_single_test_job( description="Test job", size=1, failure_rate=0, + commit_within_job=False, + failure_retry_seconds=0, ): delayed = ( http.request.env["queue.job"] @@ -259,7 +337,11 @@ def _create_single_test_job( channel=channel, description=description, ) - ._test_job(failure_rate=failure_rate) + ._test_job( + failure_rate=failure_rate, + commit_within_job=commit_within_job, + failure_retry_seconds=failure_retry_seconds, + ) ) return "job uuid: %s" % (delayed.db_record().uuid,) @@ -274,6 +356,8 @@ def _create_graph_test_jobs( channel=None, description="Test job", failure_rate=0, + commit_within_job=False, + failure_retry_seconds=0, ): model = http.request.env["queue.job"] current_count = 0 @@ -296,7 +380,11 @@ def _create_graph_test_jobs( max_retries=max_retries, channel=channel, description="%s #%d" % (description, current_count), - )._test_job(failure_rate=failure_rate) + )._test_job( + failure_rate=failure_rate, + commit_within_job=commit_within_job, + failure_retry_seconds=failure_retry_seconds, + ) ) grouping = random.choice(possible_grouping_methods) diff --git a/queue_job/data/queue_data.xml b/queue_job/data/queue_data.xml index ca5a747746..a2680cc475 100644 --- a/queue_job/data/queue_data.xml +++ b/queue_job/data/queue_data.xml @@ -1,15 +1,6 @@ - - Jobs Garbage Collector - 5 - minutes - -1 - - code - model.requeue_stuck_jobs() - Job failed diff --git a/queue_job/job.py b/queue_job/job.py index e486d1f001..86f43981a0 100644 --- a/queue_job/job.py +++ b/queue_job/job.py @@ -13,9 +13,39 @@ from random import randint import odoo +from odoo.models import BaseModel from .exception import FailedJobError, NoSuchJobError, RetryableJobError +try: + from contextlib import contextmanager, nullcontext +except ImportError: # Python 3.6 + from contextlib import contextmanager + + @contextmanager + def nullcontext(): + yield + + +def _rebind_to_cr(value, cr): + """Rebind any BaseModel inside ``value`` to the cursor ``cr``. + + Recurses into lists, tuples and dicts. Preserves uid/su/context of each + inner env - only the cursor is swapped. Recordsets already bound to + ``cr`` and non-recordset values pass through untouched. Containers are + rebuilt, so in-place changes to the result are lost. + """ + if isinstance(value, BaseModel): + if value.env.cr is cr: + return value + return value.with_env(value.env(cr=cr)) + if isinstance(value, (list, tuple)): + return type(value)(_rebind_to_cr(v, cr) for v in value) + if isinstance(value, dict): + return {k: _rebind_to_cr(v, cr) for k, v in value.items()} + return value + + WAIT_DEPENDENCIES = "wait_dependencies" PENDING = "pending" ENQUEUED = "enqueued" @@ -238,6 +268,56 @@ def load_many(cls, env, job_uuids): recordset = cls.db_records_from_uuids(env, job_uuids) return {cls._load_from_db_record(record) for record in recordset} + def add_lock_record(self) -> None: + """ + Create row in db to be locked while the job is being performed. + """ + self.env.cr.execute( + """ + INSERT INTO + queue_job_lock (id, queue_job_id) + SELECT + id, id + FROM + queue_job + WHERE + uuid = %s + ON CONFLICT(id) + DO NOTHING; + """, + [self.uuid], + ) + + def lock(self) -> bool: + """Lock row of job that is being performed. + + Return False if a job cannot be locked: it means that the job is not in + STARTED state or is already locked by another worker. + """ + self.env.cr.execute( + """ + SELECT + * + FROM + queue_job_lock + WHERE + queue_job_id in ( + SELECT + id + FROM + queue_job + WHERE + uuid = %s + AND state = %s + ) + FOR NO KEY UPDATE SKIP LOCKED; + """, + [self.uuid, STARTED], + ) + + # 1 job should be locked + return bool(self.env.cr.fetchall()) + @classmethod def _load_from_db_record(cls, job_db_record): stored = job_db_record @@ -437,17 +517,18 @@ def __init__( raise TypeError("Job accepts only methods of Models") recordset = func.__self__ - env = recordset.env self.method_name = func.__name__ self.recordset = recordset - self.env = env - self.job_model = self.env["queue.job"] - self.job_model_name = "queue.job" - self.job_config = ( self.env["queue.job.function"].sudo().job_config(self.job_function_name) ) + on_fail_method_name = self.job_config.on_fail_method_name + if on_fail_method_name and not _is_model_method( + getattr(self.recordset, on_fail_method_name, None) + ): + raise TypeError("Job accepts only methods of Models") + self.on_fail_method_name = on_fail_method_name self.state = PENDING @@ -494,10 +575,10 @@ def __init__( self.exc_message = None self.exc_info = None - if "company_id" in env.context: - company_id = env.context["company_id"] + if "company_id" in self.env.context: + company_id = self.env.context["company_id"] else: - company_id = env.company.id + company_id = self.env.company.id self.company_id = company_id self._eta = None self.eta = eta @@ -522,7 +603,12 @@ def perform(self): """ self.retry += 1 try: - self.result = self.func(*tuple(self.args), **self.kwargs) + if self.job_config.allow_commit: + env_context_manager = self.in_temporary_env() + else: + env_context_manager = nullcontext() + with env_context_manager: + self.result = self.func(*tuple(self.args), **self.kwargs) except RetryableJobError as err: if err.ignore_retry: self.retry -= 1 @@ -542,6 +628,16 @@ def perform(self): return self.result + @contextmanager + def in_temporary_env(self): + with self.env.registry.cursor() as new_cr: + env = self.env + try: + self._env = env(cr=new_cr) + yield + finally: + self._env = env + def _get_common_dependent_jobs_query(self): return """ UPDATE queue_job @@ -572,6 +668,9 @@ def _get_common_dependent_jobs_query(self): AND state = %s; """ + def should_check_dependents(self): + return any(self.__reverse_depends_on_uuids) + def enqueue_waiting(self): sql = self._get_common_dependent_jobs_query() self.env.cr.execute(sql, (PENDING, self.uuid, DONE, WAIT_DEPENDENCIES)) @@ -710,6 +809,32 @@ def __lt__(self, other): def db_record(self): return self.db_records_from_uuids(self.env, [self.uuid]) + @property + def env(self): + return self.recordset.env + + @env.setter + def _env(self, env): + self.recordset = self.recordset.with_env(env) + + @property + def args(self): + """Positional arguments, rebound to the job's current cursor.""" + return _rebind_to_cr(self._args, self.env.cr) + + @args.setter + def args(self, value): + self._args = value + + @property + def kwargs(self): + """Keyword arguments, rebound to the job's current cursor.""" + return _rebind_to_cr(self._kwargs, self.env.cr) + + @kwargs.setter + def kwargs(self, value): + self._kwargs = value + @property def func(self): recordset = self.recordset.with_context(job_uuid=self.uuid) @@ -774,7 +899,7 @@ def model_name(self): @property def user_id(self): - return self.recordset.env.uid + return self.env.uid @property def eta(self): @@ -830,6 +955,7 @@ def set_started(self): self.state = STARTED self.date_started = datetime.now() self.worker_pid = os.getpid() + self.add_lock_record() def set_done(self, result=None): self.state = DONE @@ -851,6 +977,13 @@ def set_failed(self, **kw): if v is not None: setattr(self, k, v) + def on_fail(self, fail_vals): + if not self.on_fail_method_name: + return + on_fail_func = getattr(self.recordset, self.on_fail_method_name, None) + if on_fail_func: + on_fail_func(**fail_vals) + def __repr__(self): return "" % (self.uuid, self.priority) diff --git a/queue_job/jobrunner/channels.py b/queue_job/jobrunner/channels.py index 54188db11f..d5d28d8b15 100644 --- a/queue_job/jobrunner/channels.py +++ b/queue_job/jobrunner/channels.py @@ -402,9 +402,22 @@ class Channel: with a capacity of 1. It is also possible to dedicate a channel with a limited capacity for application-autocreated subchannels without risking to overflow the system. + + A paused channel does not process any job until it is resumed. All subchannels + are blocked with their parent channel. """ - def __init__(self, name, parent, capacity=None, sequential=False, throttle=0): + def __init__( + self, + name, + parent, + capacity=None, + sequential=False, + throttle=0, + paused=False, + capacity_default=None, + sequential_default=False, + ): self.name = name self.parent = parent if self.parent: @@ -414,9 +427,14 @@ def __init__(self, name, parent, capacity=None, sequential=False, throttle=0): self._running = set() self._failed = set() self._pause_until = 0 # utc seconds since the epoch - self.capacity = capacity + self.capacity = ( + capacity if (capacity is not None) else (parent and parent.capacity_default) + ) + self.capacity_default = capacity_default self.throttle = throttle # seconds - self.sequential = sequential + self.sequential = sequential or (parent and parent.sequential_default) + self.sequential_default = sequential_default + self.paused = paused @property def sequential(self): @@ -432,13 +450,19 @@ def configure(self, config): Supported keys are: * capacity + * capacity_default (default for sub channels) * sequential + * sequential_default (default for sub channels) * throttle + * paused """ assert self.fullname.endswith(config["name"]) self.capacity = config.get("capacity", None) + self.capacity_default = config.get("capacity_default", None) self.sequential = bool(config.get("sequential", False)) + self.sequential_default = config.get("sequential_default", False) self.throttle = int(config.get("throttle", 0)) + self.paused = bool(config.get("paused", False)) if self.sequential and self.capacity != 1: raise ValueError("A sequential channel must have a capacity of 1") @@ -455,12 +479,13 @@ def get_subchannel_by_name(self, subchannel_name): def __str__(self): capacity = "∞" if self.capacity is None else str(self.capacity) - return "%s(C:%s,Q:%d,R:%d,F:%d)" % ( + return "%s(C:%s,Q:%d,R:%d,F:%d%s)" % ( self.fullname, capacity, len(self._queue), len(self._running), len(self._failed), + ",paused" if self.paused else "", ) def remove(self, job): @@ -517,6 +542,8 @@ def set_failed(self, job): _logger.debug("job %s marked failed in channel %s", job.uuid, self) def has_capacity(self): + if self.paused: + return False if self.sequential and self._failed: # a sequential queue blocks on failed jobs return False @@ -799,6 +826,31 @@ class ChannelManager: >>> cm.notify(db, 'S', 'S3', 3, 0, 10, None, 'done') >>> pp(list(cm.get_jobs_to_run(now=105))) [] + + Test pausing a channel + + >>> cm = ChannelManager() + >>> cm.simple_configure('root:4,P:2:paused,P.sub:1') + >>> cm.notify(db, 'P', 'P1', 1, 0, 10, None, 'pending') + >>> cm.notify(db, 'P.sub', 'PS1', 2, 0, 10, None, 'pending') + + Paused channel yields no job + + >>> pp(list(cm.get_jobs_to_run(now=100))) + [] + + Resuming the channel yields the pending jobs + + >>> cm.simple_configure('root:4,P:2') + >>> pp(list(cm.get_jobs_to_run(now=100))) + [, ] + + Pausing the root channel blocks everything + + >>> cm.simple_configure('root:4:paused') + >>> cm.notify(db, 'P', 'P3', 4, 0, 10, None, 'pending') + >>> pp(list(cm.get_jobs_to_run(now=106))) + [] """ def __init__(self): @@ -866,22 +918,24 @@ def parse_simple_config(cls, config_string): continue config = {} config_items = split_strip(channel_config_string, ":") - name = config_items[0] + name = config_items.pop(0) if not name: raise ValueError( "Invalid channel config %s: missing channel name" % config_string ) config["name"] = name - if len(config_items) > 1: - capacity = config_items[1] + if len(config_items) > 0: try: - config["capacity"] = int(capacity) + config["capacity"] = int(config_items[0]) + config_items.pop(0) except Exception: - raise ValueError( - "Invalid channel config %s: " - "invalid capacity %s" % (config_string, capacity) - ) - for config_item in config_items[2:]: + if name == "root": + raise ValueError( + "Invalid channel config %s: " + "invalid capacity %s" % (config_string, config_items[0]) + ) + + for config_item in config_items: kv = split_strip(config_item, "=") if len(kv) == 1: k, v = kv[0], True @@ -894,10 +948,19 @@ def parse_simple_config(cls, config_string): ) if k in config: raise ValueError( - "Invalid channel config %s: " - "duplicate key %s" % (config_string, k) + "Invalid channel config %s: duplicate key %s" + % (config_string, k) ) - config[k] = v + if k == "capacity_default": + try: + config[k] = int(v) + except Exception: + raise ValueError( + "Invalid channel config %s: " + "invalid capacity_default %s" % (config_string, v) + ) + else: + config[k] = v else: config["capacity"] = 1 res.append(config) @@ -910,6 +973,17 @@ def simple_configure(self, config_string): >>> c = cm.get_channel_by_name('root') >>> c.capacity 1 + + >>> cm.simple_configure('root:bogus') + Traceback (most recent call last): + ... + ValueError: Invalid channel config root:bogus: invalid capacity bogus + + >>> cm.simple_configure('root:4,:2') + Traceback (most recent call last): + ... + ValueError: Invalid channel config root:4,:2: missing channel name + >>> cm.simple_configure('root:4,autosub.sub:2,seq:1:sequential') >>> cm.get_channel_by_name('root').capacity 4 @@ -926,7 +1000,28 @@ def simple_configure(self, config_string): 1 >>> cm.get_channel_by_name('seq').sequential True - """ + + >>> cm.simple_configure('root:4:capacity_default=bogus') + Traceback (most recent call last): + ... + ValueError: Invalid channel config root:4:capacity_default=bogus: invalid capacity_default bogus + + >>> cm.simple_configure('root:4,sub:3:capacity_default=2') + >>> cm.get_channel_by_name('root.sub').capacity + 3 + >>> cm.get_channel_by_name('root.sub.auto', autocreate=True).capacity + 2 + + >>> cm.simple_configure('root:4,seq:2:sequential') + Traceback (most recent call last): + ... + ValueError: A sequential channel must have a capacity of 1 + + >>> cm.simple_configure('root:4,seq:sequential_default') + >>> cm.get_channel_by_name('root.seq.auto', autocreate=True).sequential + True + + """ # noqa: E501,B950 for config in ChannelManager.parse_simple_config(config_string): self.get_channel_from_config(config) diff --git a/queue_job/jobrunner/runner.py b/queue_job/jobrunner/runner.py index d4dcabceb2..bf09bfc863 100644 --- a/queue_job/jobrunner/runner.py +++ b/queue_job/jobrunner/runner.py @@ -114,22 +114,6 @@ * After creating a new database or installing queue_job on an existing database, Odoo must be restarted for the runner to detect it. -* When Odoo shuts down normally, it waits for running jobs to finish. - However, when the Odoo server crashes or is otherwise force-stopped, - running jobs are interrupted while the runner has no chance to know - they have been aborted. In such situations, jobs may remain in - ``started`` or ``enqueued`` state after the Odoo server is halted. - Since the runner has no way to know if they are actually running or - not, and does not know for sure if it is safe to restart the jobs, - it does not attempt to restart them automatically. Such stale jobs - therefore fill the running queue and prevent other jobs to start. - You must therefore requeue them manually, either from the Jobs view, - or by running the following SQL statement *before starting Odoo*: - -.. code-block:: sql - - update queue_job set state='pending' where state in ('started', 'enqueued') - .. rubric:: Footnotes .. [1] From a security standpoint, it is safe to have an anonymous HTTP @@ -139,7 +123,6 @@ of running Odoo is obviously not for production purposes. """ -import datetime import logging import os import selectors @@ -155,16 +138,21 @@ from odoo.tools import config from . import queue_job_config -from .channels import ENQUEUED, NOT_DONE, PENDING, ChannelManager +from .channels import ENQUEUED, NOT_DONE, ChannelManager SELECT_TIMEOUT = 60 ERROR_RECOVERY_DELAY = 5 +PG_ADVISORY_LOCK_ID = 2293787760715711918 _logger = logging.getLogger(__name__) select = selectors.DefaultSelector +class MasterElectionLost(Exception): + pass + + # Unfortunately, it is not possible to extend the Odoo # server command line arguments, so we resort to environment variables # to configure the runner (channels mostly). @@ -181,15 +169,10 @@ def _channels(): ) -def _datetime_to_epoch(dt): +def _odoo_now(): # important: this must return the same as postgresql # EXTRACT(EPOCH FROM TIMESTAMP dt) - return (dt - datetime.datetime(1970, 1, 1)).total_seconds() - - -def _odoo_now(): - dt = datetime.datetime.utcnow() - return _datetime_to_epoch(dt) + return time.time() def _connection_info_for(db_name): @@ -207,28 +190,6 @@ def _connection_info_for(db_name): def _async_http_get(scheme, host, port, user, password, db_name, job_uuid): - # Method to set failed job (due to timeout, etc) as pending, - # to avoid keeping it as enqueued. - def set_job_pending(): - connection_info = _connection_info_for(db_name) - conn = psycopg2.connect(**connection_info) - conn.set_isolation_level(ISOLATION_LEVEL_AUTOCOMMIT) - with closing(conn.cursor()) as cr: - cr.execute( - "UPDATE queue_job SET state=%s, " - "date_enqueued=NULL, date_started=NULL " - "WHERE uuid=%s and state=%s " - "RETURNING uuid", - (PENDING, job_uuid, ENQUEUED), - ) - if cr.fetchone(): - _logger.warning( - "state of job %s was reset from %s to %s", - job_uuid, - ENQUEUED, - PENDING, - ) - # TODO: better way to HTTP GET asynchronously (grequest, ...)? # if this was python3 I would be doing this with # asyncio, aiohttp and aiopg @@ -236,6 +197,7 @@ def urlopen(): url = "{}://{}:{}/queue_job/runjob?db={}&job_uuid={}".format( scheme, host, port, db_name, job_uuid ) + # pylint: disable=except-pass try: auth = None if user: @@ -249,10 +211,10 @@ def urlopen(): # for codes between 500 and 600 response.raise_for_status() except requests.Timeout: - set_job_pending() + # A timeout is a normal behaviour, it shouldn't be logged as an exception + pass except Exception: _logger.exception("exception in GET %s", url) - set_job_pending() thread = threading.Thread(target=urlopen) thread.daemon = True @@ -264,10 +226,15 @@ def __init__(self, db_name): self.db_name = db_name connection_info = _connection_info_for(db_name) self.conn = psycopg2.connect(**connection_info) - self.conn.set_isolation_level(ISOLATION_LEVEL_AUTOCOMMIT) - self.has_queue_job = self._has_queue_job() - if self.has_queue_job: - self._initialize() + try: + self.conn.set_isolation_level(ISOLATION_LEVEL_AUTOCOMMIT) + self.has_queue_job = self._has_queue_job() + if self.has_queue_job: + self._acquire_master_lock() + self._initialize() + except BaseException: + self.close() + raise def close(self): # pylint: disable=except-pass @@ -280,6 +247,14 @@ def close(self): pass self.conn = None + def _acquire_master_lock(self): + """Acquire the master runner lock or raise MasterElectionLost""" + with closing(self.conn.cursor()) as cr: + cr.execute("SELECT pg_try_advisory_lock(%s)", (PG_ADVISORY_LOCK_ID,)) + if not cr.fetchone()[0]: + msg = f"could not acquire master runner lock on {self.db_name}" + raise MasterElectionLost(msg) + def _has_queue_job(self): with closing(self.conn.cursor()) as cr: cr.execute( @@ -343,6 +318,105 @@ def set_job_enqueued(self, uuid): (ENQUEUED, uuid), ) + def _query_requeue_dead_jobs(self): + return """ + UPDATE + queue_job + SET + state=( + CASE + WHEN + max_retries IS NOT NULL AND + max_retries != 0 AND -- infinite retries if max_retries is 0 + retry IS NOT NULL AND + retry>max_retries + THEN 'failed' + ELSE 'pending' + END), + retry=( + CASE + WHEN state='started' + THEN COALESCE(retry,0)+1 ELSE retry + END), + exc_name=( + CASE + WHEN + max_retries IS NOT NULL AND + max_retries != 0 AND -- infinite retries if max_retries is 0 + retry IS NOT NULL AND + retry>max_retries + THEN 'JobFoundDead' + ELSE exc_name + END), + exc_info=( + CASE + WHEN + max_retries IS NOT NULL AND + max_retries != 0 AND -- infinite retries if max_retries is 0 + retry IS NOT NULL AND + retry>max_retries + THEN 'Job found dead after too many retries' + ELSE exc_info + END) + WHERE + state IN ('enqueued','started') + AND date_enqueued < (now() AT TIME ZONE 'utc' - INTERVAL '10 sec') + AND ( + id in ( + SELECT + queue_job_id + FROM + queue_job_lock + WHERE + queue_job_lock.queue_job_id = queue_job.id + FOR NO KEY UPDATE SKIP LOCKED + ) + OR NOT EXISTS ( + SELECT + 1 + FROM + queue_job_lock + WHERE + queue_job_lock.queue_job_id = queue_job.id + ) + ) + RETURNING uuid + """ + + def requeue_dead_jobs(self): + """ + Set started and enqueued jobs but not locked to pending + + A job is locked when it's being executed + When a job is killed, it releases the lock + + If the number of retries exceeds the number of max retries, + the job is set as 'failed' with the error 'JobFoundDead'. + + Adding a buffer on 'date_enqueued' to check + that it has been enqueued for more than 10sec. + This prevents from requeuing jobs before they are actually started. + + When Odoo shuts down normally, it waits for running jobs to finish. + However, when the Odoo server crashes or is otherwise force-stopped, + running jobs are interrupted while the runner has no chance to know + they have been aborted. + + This also handles orphaned jobs (enqueued but never started, no lock). + This edge case occurs when the runner marks a job as 'enqueued' + but the HTTP request to start the job never reaches the Odoo server + (e.g., due to server shutdown/crash between setting enqueued and + the controller receiving the request). + """ + + with closing(self.conn.cursor()) as cr: + query = self._query_requeue_dead_jobs() + + cr.execute(query) + + for (uuid,) in cr.fetchall(): + _logger.warning("Re-queued dead job with uuid: %s", uuid) + class QueueJobRunner: def __init__( @@ -425,7 +499,8 @@ def close_databases(self, remove_jobs=True): self.db_by_name = {} def initialize_databases(self): - for db_name in self.get_db_names(): + for db_name in sorted(self.get_db_names()): + # sorting is important to avoid deadlocks in acquiring the master lock db = Database(db_name) if db.has_queue_job: self.db_by_name[db_name] = db @@ -433,6 +508,13 @@ def initialize_databases(self): for job_data in cr: self.channel_manager.notify(db_name, *job_data) _logger.info("queue job runner ready for db %s", db_name) + else: + db.close() + + def requeue_dead_jobs(self): + for db in self.db_by_name.values(): + if db.has_queue_job: + db.requeue_dead_jobs() def run_jobs(self): now = _odoo_now() @@ -519,13 +601,14 @@ def run(self): while not self._stop: # outer loop does exception recovery try: - _logger.info("initializing database connections") + _logger.debug("initializing database connections") # TODO: how to detect new databases or databases # on which queue_job is installed after server start? self.initialize_databases() _logger.info("database connections ready") # inner loop does the normal processing while not self._stop: + self.requeue_dead_jobs() self.process_notifications() self.run_jobs() self.wait_notification() @@ -534,6 +617,14 @@ def run(self): except InterruptedError: # Interrupted system call, i.e. KeyboardInterrupt during select self.stop() + except MasterElectionLost as e: + _logger.debug( + "master election lost: %s, sleeping %ds and retrying", + e, + ERROR_RECOVERY_DELAY, + ) + self.close_databases() + time.sleep(ERROR_RECOVERY_DELAY) except Exception: _logger.exception( "exception: sleeping %ds and retrying", ERROR_RECOVERY_DELAY diff --git a/queue_job/migrations/14.0.3.16.0/pre-migration.py b/queue_job/migrations/14.0.3.16.0/pre-migration.py new file mode 100644 index 0000000000..931c336866 --- /dev/null +++ b/queue_job/migrations/14.0.3.16.0/pre-migration.py @@ -0,0 +1,11 @@ +# License LGPL-3.0 or later (http://www.gnu.org/licenses/lgpl.html) +from openupgradelib import openupgrade + + +@openupgrade.migrate() +def migrate(env, version): + # Remove cron garbage collector + openupgrade.delete_records_safely_by_xml_id( + env, + ["queue_job.ir_cron_queue_job_garbage_collector"], + ) diff --git a/queue_job/models/__init__.py b/queue_job/models/__init__.py index 4744e7ab46..6265dfe9cb 100644 --- a/queue_job/models/__init__.py +++ b/queue_job/models/__init__.py @@ -3,3 +3,4 @@ from . import queue_job from . import queue_job_channel from . import queue_job_function +from . import queue_job_lock diff --git a/queue_job/models/queue_job.py b/queue_job/models/queue_job.py index 5c69655181..13fbc08b2d 100644 --- a/queue_job/models/queue_job.py +++ b/queue_job/models/queue_job.py @@ -6,17 +6,17 @@ from datetime import datetime, timedelta from odoo import _, api, exceptions, fields, models -from odoo.osv import expression from odoo.tools import config, html_escape, index_exists from odoo.addons.base_sparse_field.models.fields import Serialized from ..delay import Graph -from ..exception import JobError +from ..exception import JobError, RetryableJobError from ..fields import JobSerialized from ..job import ( CANCELLED, DONE, + ENQUEUED, FAILED, PENDING, STARTED, @@ -104,6 +104,7 @@ class QueueJob(models.Model): date_done = fields.Datetime(readonly=True) exec_time = fields.Float( string="Execution Time (avg)", + readonly=True, group_operator="avg", help="Time required to execute this job in seconds. Average when grouped.", ) @@ -327,18 +328,26 @@ def _change_job_state(self, state, result=None): raise ValueError("State not supported: %s" % state) def button_done(self): + # If job was set to STARTED or CANCELLED, do not set it to DONE + states_from = (WAIT_DEPENDENCIES, PENDING, ENQUEUED, FAILED) result = _("Manually set to done by %s") % self.env.user.name - self._change_job_state(DONE, result=result) + records = self.filtered(lambda job_: job_.state in states_from) + records._change_job_state(DONE, result=result) return True def button_cancelled(self): + # If job was set to DONE do not cancel it + states_from = (WAIT_DEPENDENCIES, PENDING, ENQUEUED, FAILED) result = _("Cancelled by %s") % self.env.user.name - self._change_job_state(CANCELLED, result=result) + records = self.filtered(lambda job_: job_.state in states_from) + records._change_job_state(CANCELLED, result=result) return True def requeue(self): - jobs_to_requeue = self.filtered(lambda job_: job_.state != WAIT_DEPENDENCIES) - jobs_to_requeue._change_job_state(PENDING) + # If job is already in queue or started, do not requeue it + states_from = (FAILED, DONE, CANCELLED) + records = self.filtered(lambda job_: job_.state in states_from) + records._change_job_state(PENDING) return True def _message_post_on_failure(self): @@ -346,8 +355,11 @@ def _message_post_on_failure(self): # at every job creation domain = self._subscribe_users_domain() base_users = self.env["res.users"].search(domain) + suscribe_job_creator = self._subscribe_job_creator() for record in self: - users = base_users | record.user_id + users = base_users + if suscribe_job_creator: + users |= record.user_id record.message_subscribe(partner_ids=users.mapped("partner_id").ids) msg = record._message_failed_job() if msg: @@ -364,6 +376,14 @@ def _subscribe_users_domain(self): domain.append(("company_id", "in", companies.ids)) return domain + @api.model + def _subscribe_job_creator(self): + """ + Whether the user that created the job should be subscribed to the job, + in addition to users determined by `_subscribe_users_domain` + """ + return True + def _message_failed_job(self): """Return a message which will be posted on the job when it is failed. @@ -405,62 +425,13 @@ def autovacuum(self): limit=1000, ) if jobs: - jobs.unlink() + jobs.sudo().unlink() if not config["test_enable"]: self.env.cr.commit() # pylint: disable=E8102 else: break return True - def requeue_stuck_jobs(self, enqueued_delta=5, started_delta=0): - """Fix jobs that are in a bad states - - :param in_queue_delta: lookup time in minutes for jobs - that are in enqueued state - - :param started_delta: lookup time in minutes for jobs - that are in enqueued state, - 0 means that it is not checked - """ - self._get_stuck_jobs_to_requeue( - enqueued_delta=enqueued_delta, started_delta=started_delta - ).requeue() - return True - - def _get_stuck_jobs_domain(self, queue_dl, started_dl): - domain = [] - now = fields.datetime.now() - if queue_dl: - queue_dl = now - timedelta(minutes=queue_dl) - domain.append( - [ - "&", - ("date_enqueued", "<=", fields.Datetime.to_string(queue_dl)), - ("state", "=", "enqueued"), - ] - ) - if started_dl: - started_dl = now - timedelta(minutes=started_dl) - domain.append( - [ - "&", - ("date_started", "<=", fields.Datetime.to_string(started_dl)), - ("state", "=", "started"), - ] - ) - if not domain: - raise exceptions.ValidationError( - _("If both parameters are 0, ALL jobs will be requeued!") - ) - return expression.OR(domain) - - def _get_stuck_jobs_to_requeue(self, enqueued_delta, started_delta): - job_model = self.env["queue.job"] - stuck_jobs = job_model.search( - self._get_stuck_jobs_domain(enqueued_delta, started_delta) - ) - return stuck_jobs - def related_action_open_record(self): """Open a form view with the record(s) of the job. @@ -494,7 +465,24 @@ def related_action_open_record(self): ) return action - def _test_job(self, failure_rate=0): + def _test_job( + self, + failure_rate=0, + commit_within_job=False, + failure_retry_seconds=0, + ): _logger.info("Running test job.") if random.random() <= failure_rate: - raise JobError("Job failed") + if failure_retry_seconds: + raise RetryableJobError( + f"Retryable job failed, will be retried in " + f"{failure_retry_seconds} seconds", + seconds=failure_retry_seconds, + ) + else: + raise JobError("Job failed") + if commit_within_job: + self.env.cr.commit() # pylint: disable=invalid-commit + + def _test_on_fail(self, **kw): + pass diff --git a/queue_job/models/queue_job_function.py b/queue_job/models/queue_job_function.py index d839708dcf..af9b1f6e5b 100644 --- a/queue_job/models/queue_job_function.py +++ b/queue_job/models/queue_job_function.py @@ -28,7 +28,9 @@ class QueueJobFunction(models.Model): "related_action_enable " "related_action_func_name " "related_action_kwargs " - "job_function_id ", + "job_function_id " + "allow_commit " + "on_fail_method_name", ) def _default_channel(self): @@ -47,6 +49,10 @@ def _default_channel(self): comodel_name="ir.model", string="Model", ondelete="cascade" ) method = fields.Char() + on_fail_method = fields.Char( + help="Model function to be called if the job is failed and will not be " + "retried.", + ) channel_id = fields.Many2one( comodel_name="queue.job.channel", @@ -79,6 +85,12 @@ def _default_channel(self): "enable, func_name, kwargs.\n" "See the module description for details.", ) + allow_commit = fields.Boolean( + help="Allows the job to commit transactions during execution. " + "Under the hood, this executes the job in a new database cursor, " + "which incurs an overhead as it requires an extra connection to " + "the database. " + ) @api.depends("model_id.model", "method") def _compute_name(self): @@ -143,6 +155,8 @@ def job_default_config(self): related_action_func_name=None, related_action_kwargs={}, job_function_id=None, + allow_commit=False, + on_fail_method_name=None, ) def _parse_retry_pattern(self): @@ -178,6 +192,8 @@ def job_config(self, name): related_action_func_name=config.related_action.get("func_name"), related_action_kwargs=config.related_action.get("kwargs", {}), job_function_id=config.id, + allow_commit=config.allow_commit, + on_fail_method_name=config.on_fail_method, ) def _retry_pattern_format_error_message(self): diff --git a/queue_job/models/queue_job_lock.py b/queue_job/models/queue_job_lock.py new file mode 100644 index 0000000000..b01c7f3a91 --- /dev/null +++ b/queue_job/models/queue_job_lock.py @@ -0,0 +1,16 @@ +# Copyright 2025 ACSONE SA/NV +# License AGPL-3.0 or later (https://www.gnu.org/licenses/agpl). + +from odoo import fields, models + + +class QueueJobLock(models.Model): + _name = "queue.job.lock" + _description = "Queue Job Lock" + + queue_job_id = fields.Many2one( + comodel_name="queue.job", + required=True, + ondelete="cascade", + index=True, + ) diff --git a/queue_job/security/ir.model.access.csv b/queue_job/security/ir.model.access.csv index 634daf8ede..9a451d6a78 100644 --- a/queue_job/security/ir.model.access.csv +++ b/queue_job/security/ir.model.access.csv @@ -1,5 +1,6 @@ id,name,model_id:id,group_id:id,perm_read,perm_write,perm_create,perm_unlink -access_queue_job_manager,queue job manager,queue_job.model_queue_job,queue_job.group_queue_job_manager,1,1,1,1 +access_queue_job_manager,queue job manager,queue_job.model_queue_job,queue_job.group_queue_job_manager,1,1,0,0 +access_queue_job_lock_manager,queue job lock manager,queue_job.model_queue_job_lock,queue_job.group_queue_job_manager,1,0,0,0 access_queue_job_function_manager,queue job functions manager,queue_job.model_queue_job_function,queue_job.group_queue_job_manager,1,1,1,1 access_queue_job_channel_manager,queue job channel manager,queue_job.model_queue_job_channel,queue_job.group_queue_job_manager,1,1,1,1 access_queue_requeue_job,queue requeue job manager,queue_job.model_queue_requeue_job,queue_job.group_queue_job_manager,1,1,1,1 diff --git a/queue_job/tests/__init__.py b/queue_job/tests/__init__.py index db53ac3a60..51fd9816ba 100644 --- a/queue_job/tests/__init__.py +++ b/queue_job/tests/__init__.py @@ -1,3 +1,4 @@ +from . import test_run_job_controller from . import test_runner_channels from . import test_runner_runner from . import test_delayable diff --git a/queue_job/tests/common.py b/queue_job/tests/common.py index 76c03218ba..b8ab1da9f7 100644 --- a/queue_job/tests/common.py +++ b/queue_job/tests/common.py @@ -1,15 +1,11 @@ # Copyright 2019 Camptocamp # License AGPL-3.0 or later (http://www.gnu.org/licenses/agpl). -import doctest -import logging -import sys import typing from contextlib import contextmanager from itertools import groupby from operator import attrgetter from unittest import TestCase, mock -# pylint: disable=odoo-addons-relative-import from odoo.addons.queue_job.delay import Graph # pylint: disable=odoo-addons-relative-import @@ -95,7 +91,7 @@ def button_that_uses_delayable_chain(self): with mock.patch( "odoo.addons.queue_job.delay.Job", name="Job Class", - auto_spec=True, + unsafe=True, ) as job_cls_mock: with JobsTrap(job_cls_mock) as trap: yield trap @@ -277,7 +273,7 @@ def _add_job(self, *args, **kwargs): def _prepare_context(self, job): # pylint: disable=context-overridden - job_model = job.job_model.with_context({}) + job_model = job.env["queue.job"].with_context({}) field_records = job_model._fields["records"] # Filter the context to simulate store/load of the job job.recordset = field_records.convert_to_write(job.recordset, job_model) @@ -413,54 +409,3 @@ def test_export(self): delayable = mock.MagicMock(name="DelayableBinding") delayable_cls.return_value = delayable yield delayable_cls, delayable - - -class OdooDocTestCase(doctest.DocTestCase): - """ - We need a custom DocTestCase class in order to: - - define test_tags to run as part of standard tests - - output a more meaningful test name than default "DocTestCase.runTest" - """ - - def __init__( - self, doctest, optionflags=0, setUp=None, tearDown=None, checker=None, seq=0 - ): - super().__init__( - doctest._dt_test, - optionflags=optionflags, - setUp=setUp, - tearDown=tearDown, - checker=checker, - ) - self.test_sequence = seq - - def setUp(self): - """Log an extra statement which test is started.""" - super().setUp() - logging.getLogger(__name__).info("Running tests for %s", self._dt_test.name) - - -def load_doctests(module): - """ - Generates a tests loading method for the doctests of the given module - https://docs.python.org/3/library/unittest.html#load-tests-protocol - """ - - def load_tests(loader, tests, ignore): - """ - Apply the 'test_tags' attribute to each DocTestCase found by the DocTestSuite. - Also extend the DocTestCase class trivially to fit the class teardown - that Odoo backported for its own test classes from Python 3.8. - """ - if sys.version_info < (3, 8): - doctest.DocTestCase.doClassCleanups = lambda: None - doctest.DocTestCase.tearDown_exceptions = [] - - for idx, test in enumerate(doctest.DocTestSuite(module)): - odoo_test = OdooDocTestCase(test, seq=idx) - odoo_test.test_tags = {"standard", "at_install", "queue_job", "doctest"} - tests.addTest(odoo_test) - - return tests - - return load_tests diff --git a/queue_job/tests/test_json_field.py b/queue_job/tests/test_json_field.py index 802911c9eb..7d2533df93 100644 --- a/queue_job/tests/test_json_field.py +++ b/queue_job/tests/test_json_field.py @@ -20,16 +20,20 @@ def test_encoder_recordset(self): partner = self.env(user=demo_user, context=context).ref("base.main_partner") value = partner value_json = json.dumps(value, cls=JobEncoder) + expected_context = context.copy() expected = { "uid": demo_user.id, "_type": "odoo_recordset", "model": "res.partner", "ids": [partner.id], "su": False, - # no allowed context by default, must be changed in 16.0 - "context": {}, } - self.assertEqual(json.loads(value_json), expected) + result_dict = json.loads(value_json) + result_context = result_dict.pop("context") + self.assertEqual(result_dict, expected) + # context is tested separately as the order/amount of keys is not guaranteed + for key in result_context: + self.assertEqual(result_context[key], expected_context[key]) def test_encoder_recordset_list(self): demo_user = self.env.ref("base.user_demo") @@ -50,7 +54,20 @@ def test_encoder_recordset_list(self): "context": {}, }, ] - self.assertEqual(json.loads(value_json), expected) + result_dict = json.loads(value_json) + for result_value, expected_value in zip(result_dict, expected): + if isinstance(expected_value, dict): + for key in result_value: + if key == "context": + for context_key in result_value["context"]: + self.assertEqual( + result_value["context"][context_key], + expected_value["context"][context_key], + ) + else: + self.assertEqual(result_value[key], expected_value[key]) + else: + self.assertEqual(result_value, expected_value) def test_decoder_recordset(self): demo_user = self.env.ref("base.user_demo") diff --git a/queue_job/tests/test_model_job_function.py b/queue_job/tests/test_model_job_function.py index e6ddf3fcc3..dab30c78a1 100644 --- a/queue_job/tests/test_model_job_function.py +++ b/queue_job/tests/test_model_job_function.py @@ -35,6 +35,7 @@ def test_function_job_config(self): { "model_id": self.env.ref("base.model_res_users").id, "method": "read", + "on_fail_method": "search_read", "channel_id": channel.id, "edit_retry_pattern": "{1: 2, 3: 4}", "edit_related_action": ( @@ -42,6 +43,7 @@ def test_function_job_config(self): ' "func_name": "related_action_foo",' ' "kwargs": {"b": 1}}' ), + "allow_commit": True, } ) self.assertEqual( @@ -53,5 +55,7 @@ def test_function_job_config(self): related_action_func_name="related_action_foo", related_action_kwargs={"b": 1}, job_function_id=job_function.id, + allow_commit=True, + on_fail_method_name="search_read", ), ) diff --git a/queue_job/tests/test_run_job_controller.py b/queue_job/tests/test_run_job_controller.py new file mode 100644 index 0000000000..a9c89b0592 --- /dev/null +++ b/queue_job/tests/test_run_job_controller.py @@ -0,0 +1,61 @@ +# License AGPL-3.0 or later (https://www.gnu.org/licenses/agpl). +from unittest.mock import patch + +from odoo.tests.common import TransactionCase +from odoo.tools import mute_logger + +from ..controllers.main import RunJobController +from ..exception import JobError +from ..job import Job + + +class TestRunJobController(TransactionCase): + def setUp(self): + super().setUp() + + def _clean_queue_job(): + self.env["queue.job"].search([]).unlink() + + self.addCleanup(_clean_queue_job) + + def test_get_failure_values(self): + method = self.env["res.users"].mapped + job = Job(method) + ctrl = RunJobController() + rslt = ctrl._get_failure_values(job, "info", Exception("zero", "one")) + self.assertEqual( + rslt, {"exc_info": "info", "exc_name": "Exception", "exc_message": "zero"} + ) + + def test_runjob_success(self): + job = self.env["queue.job"].with_delay()._test_job() + RunJobController._runjob(self.env, job) + self.assertEqual(job.state, "done") + self.assertEqual(job.db_record().state, "done") + + def test_runjob_on_fail(self): + function = self.env.ref("queue_job.job_function_queue_job__test_job") + function.on_fail_method = "_test_on_fail" + job = self.env["queue.job"].with_delay()._test_job(failure_rate=1) + patch_test_on_fail = patch( + "odoo.addons.queue_job.models.queue_job.QueueJob._test_on_fail" + ) + patch_temp_env = patch("odoo.addons.queue_job.job.Job.in_temporary_env") + with patch_test_on_fail as mocked_hook, patch_temp_env as mocked_temp_env: + mocked_temp_env.return_value.__enter__.return_value = self.env + with self.assertRaises(JobError), mute_logger( + "odoo.addons.queue_job.controllers.main" + ): + RunJobController._runjob(self.env, job) + self.assertEqual(job.state, "failed") + self.assertEqual(mocked_hook.call_count, 1) + + def test_runjob_on_fail_not_configured(self): + job = self.env["queue.job"].with_delay()._test_job(failure_rate=1) + with patch("odoo.addons.queue_job.job.Job.in_temporary_env") as mocked_temp_env: + mocked_temp_env.return_value.__enter__.return_value = self.env + with self.assertRaises(JobError), mute_logger( + "odoo.addons.queue_job.controllers.main" + ): + RunJobController._runjob(self.env, job) + self.assertEqual(job.state, "failed") diff --git a/queue_job/tests/test_runner_channels.py b/queue_job/tests/test_runner_channels.py index d323d00683..313e350bd3 100644 --- a/queue_job/tests/test_runner_channels.py +++ b/queue_job/tests/test_runner_channels.py @@ -1,10 +1,18 @@ # Copyright 2015-2016 Camptocamp SA # License LGPL-3.0 or later (http://www.gnu.org/licenses/lgpl.html) +import doctest + +from odoo.tests import BaseCase, tagged # pylint: disable=odoo-addons-relative-import # we are testing, we want to test as we were an external consumer of the API from odoo.addons.queue_job.jobrunner import channels -from .common import load_doctests -load_tests = load_doctests(channels) +@tagged("doctest") +class TestDoctest(BaseCase): + def test_doctest(self): + results = doctest.testmod( + channels, exclude_empty=True, optionflags=doctest.REPORT_ONLY_FIRST_FAILURE + ) + self.assertEqual(results.failed, 0, "doctest failed") diff --git a/queue_job/tests/test_runner_runner.py b/queue_job/tests/test_runner_runner.py index 131ce6322d..04c3fc1946 100644 --- a/queue_job/tests/test_runner_runner.py +++ b/queue_job/tests/test_runner_runner.py @@ -1,17 +1,22 @@ # Copyright 2015-2016 Camptocamp SA # License LGPL-3.0 or later (http://www.gnu.org/licenses/lgpl.html) - -# pylint: disable=odoo-addons-relative-import -# we are testing, we want to test as we were an external consumer of the API +import doctest import os from odoo.tests import BaseCase, tagged +# pylint: disable=odoo-addons-relative-import +# we are testing, we want to test as we were an external consumer of the API from odoo.addons.queue_job.jobrunner import runner -from .common import load_doctests -load_tests = load_doctests(runner) +@tagged("doctest") +class TestDoctest(BaseCase): + def test_doctest(self): + results = doctest.testmod( + runner, exclude_empty=True, optionflags=doctest.REPORT_ONLY_FIRST_FAILURE + ) + self.assertEqual(results.failed, 0, "doctest failed") @tagged("-at_install", "post_install") diff --git a/queue_job/tests/test_wizards.py b/queue_job/tests/test_wizards.py index 2ac162d313..54718f658a 100644 --- a/queue_job/tests/test_wizards.py +++ b/queue_job/tests/test_wizards.py @@ -46,3 +46,55 @@ def test_03_done(self): wizard = self._wizard("queue.jobs.to.done") wizard.set_done() self.assertEqual(self.job.state, "done") + + def test_04_requeue_forbidden(self): + wizard = self._wizard("queue.requeue.job") + + # State WAIT_DEPENDENCIES is not requeued + self.job.state = "wait_dependencies" + wizard.requeue() + self.assertEqual(self.job.state, "wait_dependencies") + + # State PENDING, ENQUEUED or STARTED are ignored too + for test_state in ("pending", "enqueued", "started"): + self.job.state = test_state + wizard.requeue() + self.assertEqual(self.job.state, test_state) + + # States CANCELLED, DONE or FAILED will change status + self.job.state = "cancelled" + wizard.requeue() + self.assertEqual(self.job.state, "pending") + + def test_05_cancel_forbidden(self): + wizard = self._wizard("queue.jobs.to.cancelled") + + # State DONE is not cancelled + self.job.state = "done" + wizard.set_cancelled() + self.assertEqual(self.job.state, "done") + + # State PENDING, ENQUEUED, WAIT_DEPENDENCIES or FAILED will be cancelled + for test_state in ("pending", "enqueued", "wait_dependencies", "failed"): + self.job.state = test_state + wizard.set_cancelled() + self.assertEqual(self.job.state, "cancelled") + + def test_06_done_forbidden(self): + wizard = self._wizard("queue.jobs.to.done") + + # State STARTED is not set DONE manually + self.job.state = "started" + wizard.set_done() + self.assertEqual(self.job.state, "started") + + # State CANCELLED is not cancelled + self.job.state = "cancelled" + wizard.set_done() + self.assertEqual(self.job.state, "cancelled") + + # State WAIT_DEPENDENCIES, PENDING, ENQUEUED or FAILED will be set to DONE + for test_state in ("wait_dependencies", "pending", "enqueued", "failed"): + self.job.state = test_state + wizard.set_done() + self.assertEqual(self.job.state, "done") diff --git a/queue_job/utils.py b/queue_job/utils.py index 5134cd1068..16ab66e17f 100644 --- a/queue_job/utils.py +++ b/queue_job/utils.py @@ -4,6 +4,8 @@ import logging import os +from odoo import tools + _logger = logging.getLogger(__name__) @@ -36,5 +38,6 @@ def must_run_without_delay(env): return True if env.context.get("queue_job__no_delay"): - _logger.warning("`queue_job__no_delay` ctx key found. NO JOB scheduled.") + if not tools.config["test_enable"]: + _logger.info("`queue_job__no_delay` ctx key found. NO JOB scheduled.") return True diff --git a/queue_job/views/queue_job_function_views.xml b/queue_job/views/queue_job_function_views.xml index a6e2ce402c..6c208b2b67 100644 --- a/queue_job/views/queue_job_function_views.xml +++ b/queue_job/views/queue_job_function_views.xml @@ -11,6 +11,7 @@ + @@ -25,6 +26,7 @@ + diff --git a/queue_job/views/queue_job_views.xml b/queue_job/views/queue_job_views.xml index c0388c7012..abd52dc20b 100644 --- a/queue_job/views/queue_job_views.xml +++ b/queue_job/views/queue_job_views.xml @@ -25,7 +25,7 @@ />