From 0e56e44416d478e4bda1fbc6cdb0449d095abad2 Mon Sep 17 00:00:00 2001 From: Akim Juillerat Date: Fri, 24 Jul 2026 18:48:40 +0200 Subject: [PATCH] [IMP] queue_job: Implement on fail hook Allows to execute a model function when the job fails and will not be retried. [REF] queue_job: Move on fail definition to job function Rename on_fail_hook to on_fail --- queue_job/controllers/main.py | 1 + queue_job/job.py | 10 ++++++++ queue_job/models/queue_job.py | 3 +++ queue_job/models/queue_job_function.py | 9 ++++++- queue_job/tests/test_model_job_function.py | 2 ++ queue_job/tests/test_run_rob_controller.py | 28 ++++++++++++++++++++++ test_queue_job/tests/test_autovacuum.py | 1 - 7 files changed, 52 insertions(+), 2 deletions(-) diff --git a/queue_job/controllers/main.py b/queue_job/controllers/main.py index 18e257f1c6..91cda8ca3b 100644 --- a/queue_job/controllers/main.py +++ b/queue_job/controllers/main.py @@ -186,6 +186,7 @@ def retry_postpone(job, message, seconds=None): vals = cls._get_failure_values(job, traceback_txt, orig_exception) job.set_failed(**vals) job.store() + job.on_fail(vals) buff.close() raise diff --git a/queue_job/job.py b/queue_job/job.py index 032bdb9339..98d47fd050 100644 --- a/queue_job/job.py +++ b/queue_job/job.py @@ -409,6 +409,11 @@ def __init__( 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: + if 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 @@ -829,6 +834,11 @@ def set_failed(self, **kw): if v is not None: setattr(self, k, v) + def on_fail(self, fail_vals): + 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/models/queue_job.py b/queue_job/models/queue_job.py index 011edf09db..fad95fb293 100644 --- a/queue_job/models/queue_job.py +++ b/queue_job/models/queue_job.py @@ -485,3 +485,6 @@ def _test_job( time.sleep(job_duration) 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 edf90c9ab7..d29014adf4 100644 --- a/queue_job/models/queue_job_function.py +++ b/queue_job/models/queue_job_function.py @@ -29,7 +29,8 @@ class QueueJobFunction(models.Model): "related_action_func_name " "related_action_kwargs " "job_function_id " - "allow_commit", + "allow_commit " + "on_fail_method_name", ) def _default_channel(self): @@ -48,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", @@ -157,6 +162,7 @@ def job_default_config(self): related_action_kwargs={}, job_function_id=None, allow_commit=False, + on_fail_method_name=None, ) def _parse_retry_pattern(self): @@ -193,6 +199,7 @@ def job_config(self, 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/tests/test_model_job_function.py b/queue_job/tests/test_model_job_function.py index 9095f2a55e..b835c3c4ee 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": ( @@ -55,5 +56,6 @@ def test_function_job_config(self): 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_rob_controller.py b/queue_job/tests/test_run_rob_controller.py index 1a15f4363a..5e8952a94e 100644 --- a/queue_job/tests/test_run_rob_controller.py +++ b/queue_job/tests/test_run_rob_controller.py @@ -1,12 +1,23 @@ # 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(cls): + super().setUp() + + def _clean_queue_job(): + cls.env["queue.job"].search([]).unlink() + + cls.addCleanup(_clean_queue_job) + def test_get_failure_values(self): method = self.env["res.users"].mapped job = Job(method) @@ -21,3 +32,20 @@ def test_runjob_success(self): 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) + with ( + self.assertRaises(JobError), + patch( + "odoo.addons.queue_job.models.queue_job.QueueJob._test_on_fail" + ) as mocked_hook, + patch("odoo.addons.queue_job.job.Job.in_temporary_env") as mocked_temp_env, + mute_logger("odoo.addons.queue_job.controllers.main"), + ): + mocked_temp_env.return_value.__enter__.return_value = self.env + RunJobController._runjob(self.env, job) + self.assertEqual(job.state, "failed") + self.assertEqual(mocked_hook.call_count, 1) diff --git a/test_queue_job/tests/test_autovacuum.py b/test_queue_job/tests/test_autovacuum.py index 97aebcba1e..fd1cf8fedd 100644 --- a/test_queue_job/tests/test_autovacuum.py +++ b/test_queue_job/tests/test_autovacuum.py @@ -51,7 +51,6 @@ def test_autovacuum_multi_channel(self): job_60days.write( {"channel": channel_60days.complete_name, "date_done": date_done} ) - self.assertEqual( len(self.env["queue.job"].search([("channel", "!=", False)])), 2 )