From 33b3997b29d0a69b2ff25ff14d0e604d94d3b15e Mon Sep 17 00:00:00 2001 From: Florent Xicluna Date: Wed, 23 Sep 2026 07:24:23 +0200 Subject: [PATCH] [IMP] queue_job: use queue_job row lock instead of separate queue_job_lock table --- queue_job/controllers/main.py | 24 ++--- queue_job/job.py | 66 ++++---------- queue_job/jobrunner/runner.py | 89 ++++++------------- queue_job/models/__init__.py | 1 - queue_job/models/queue_job_lock.py | 16 ---- queue_job/security/ir.model.access.csv | 1 - test_queue_job/tests/common.py | 10 +++ test_queue_job/tests/test_acquire_job.py | 12 +-- test_queue_job/tests/test_requeue_dead_job.py | 42 +-------- 9 files changed, 69 insertions(+), 192 deletions(-) delete mode 100644 queue_job/models/queue_job_lock.py diff --git a/queue_job/controllers/main.py b/queue_job/controllers/main.py index 91cda8ca3b..9014a3bb87 100644 --- a/queue_job/controllers/main.py +++ b/queue_job/controllers/main.py @@ -18,7 +18,7 @@ from ..delay import chain, group from ..exception import FailedJobError, RetryableJobError -from ..job import ENQUEUED, Job +from ..job import ENQUEUED, STARTED, Job _logger = logging.getLogger(__name__) @@ -31,7 +31,7 @@ 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 + Committing while the job is not finished would release the job lock, causing it to be started again by the dead jobs requeuer. """ @@ -61,16 +61,12 @@ def _acquire_job(cls, env: api.Environment, job_uuid: str) -> Job | None: - 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. + If successful, return the Job instance, otherwise return None. + This function may fail to acquire the job, if not in the expected state + or if 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(): + job = Job.load(env, job_uuid, raise_if_not_found=False) + if not job or not job.lock(ENQUEUED): _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", @@ -78,12 +74,10 @@ def _acquire_job(cls, env: api.Environment, job_uuid: str) -> Job | None: ENQUEUED, ) return None - job = Job.load(env, job_uuid) - assert job and job.state == ENQUEUED job.set_started() job.store() env.cr.commit() - if not job.lock(): + if not job.lock(STARTED): _logger.warning( "was requested to run job %s, but it could not be locked", job_uuid, @@ -93,7 +87,7 @@ def _acquire_job(cls, env: api.Environment, job_uuid: str) -> Job | None: @classmethod def _try_perform_job(cls, env, job): - """Try to perform the job, mark it done and commit if successful.""" + """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 diff --git a/queue_job/job.py b/queue_job/job.py index a69a71ddc5..771113c018 100644 --- a/queue_job/job.py +++ b/queue_job/job.py @@ -224,15 +224,16 @@ class Job: """ @classmethod - def load(cls, env, job_uuid): + def load(cls, env, job_uuid, raise_if_not_found=True): """Read a single job from the Database Raise an error if the job is not found. """ stored = cls.db_records_from_uuids(env, [job_uuid]) - if not stored: + if stored: + return cls._load_from_db_record(stored) + if raise_if_not_found: raise NoSuchJobError(f"Job {job_uuid} does no longer exist in the storage.") - return cls._load_from_db_record(stored) @classmethod def load_many(cls, env, job_uuids): @@ -243,55 +244,23 @@ 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. + def lock(self, state) -> bool: + """Lock 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. + expected state or is already locked by another worker. + Lock is released at the commit or rollback of the transaction. """ - 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], + lock_query = ( + "SELECT uuid FROM queue_job WHERE uuid=%s AND state=%s" + " FOR NO KEY UPDATE SKIP LOCKED;" ) - - # 1 job should be locked - return bool(self.env.cr.fetchall()) + self.env.cr.execute(lock_query, [self.uuid, state]) + if not self.env.cr.fetchone(): + _logger.debug("Lock NOT acquired on %s Job %s", state, self.uuid) + return False + _logger.debug("Lock acquired on %s Job %s", state, self.uuid) + return True @classmethod def _load_from_db_record(cls, job_db_record): @@ -852,7 +821,6 @@ 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 diff --git a/queue_job/jobrunner/runner.py b/queue_job/jobrunner/runner.py index 5cbd52ef48..2d9ce8bbb6 100644 --- a/queue_job/jobrunner/runner.py +++ b/queue_job/jobrunner/runner.py @@ -354,69 +354,32 @@ def set_job_enqueued(self, 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 - """ + return """\ + WITH dead_job AS ( + SELECT id, + (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) "stop_retry", + (CASE WHEN state='started' + THEN COALESCE(retry, 0) + 1 ELSE retry END) "retry" + FROM queue_job + WHERE state IN ('enqueued', 'started') + AND date_enqueued < now() AT TIME ZONE 'utc' - INTERVAL '10 sec' + FOR NO KEY UPDATE SKIP LOCKED) + + UPDATE queue_job + SET retry = dead_job.retry, + state = CASE WHEN stop_retry THEN 'failed' ELSE 'pending' END, + exc_name = CASE WHEN stop_retry THEN 'JobFoundDead' ELSE exc_name END, + exc_info = + CASE WHEN stop_retry + THEN 'Job found dead after too many retries' + ELSE exc_info END + FROM dead_job + WHERE queue_job.id = dead_job.id + + RETURNING uuid;""" def requeue_dead_jobs(self): """ diff --git a/queue_job/models/__init__.py b/queue_job/models/__init__.py index 6265dfe9cb..4744e7ab46 100644 --- a/queue_job/models/__init__.py +++ b/queue_job/models/__init__.py @@ -3,4 +3,3 @@ 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_lock.py b/queue_job/models/queue_job_lock.py deleted file mode 100644 index b01c7f3a91..0000000000 --- a/queue_job/models/queue_job_lock.py +++ /dev/null @@ -1,16 +0,0 @@ -# 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 9a451d6a78..f9511a615e 100644 --- a/queue_job/security/ir.model.access.csv +++ b/queue_job/security/ir.model.access.csv @@ -1,6 +1,5 @@ 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,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/test_queue_job/tests/common.py b/test_queue_job/tests/common.py index d3173a2198..f7f19073dc 100644 --- a/test_queue_job/tests/common.py +++ b/test_queue_job/tests/common.py @@ -30,3 +30,13 @@ def _get_demo_job(self, uuid): "to make this test work", ) return job + + def is_job_locked(self, job, cr=None): + lock_query = ( + "SELECT 1 FROM queue_job WHERE uuid = %s FOR NO KEY UPDATE SKIP LOCKED" + ) + with self.env.registry.cursor() as cr: + cr.execute(lock_query, [job.uuid]) + if not cr.fetchone(): + return True + return False diff --git a/test_queue_job/tests/test_acquire_job.py b/test_queue_job/tests/test_acquire_job.py index 3f0c92a2be..cb8bf8b841 100644 --- a/test_queue_job/tests/test_acquire_job.py +++ b/test_queue_job/tests/test_acquire_job.py @@ -15,10 +15,8 @@ class TestRequeueDeadJob(JobCommonCase): def test_acquire_enqueued_job(self): job_record = self._get_demo_job(uuid="test_enqueued_job") self.assertFalse( - self.env["queue.job.lock"].search( - [("queue_job_id", "=", job_record.id)], - ), - "A job lock record should not exist at this point", + self.is_job_locked(job_record), + "A job lock should not exist at this point", ) with mock.patch.object( self.env.cr, "commit", mock.Mock(side_effect=self.env.flush_all) @@ -29,10 +27,8 @@ def test_acquire_enqueued_job(self): self.assertEqual(job.uuid, "test_enqueued_job") self.assertEqual(job.state, "started") self.assertTrue( - self.env["queue.job.lock"].search( - [("queue_job_id", "=", job_record.id)] - ), - "A job lock record should exist at this point", + self.is_job_locked(job_record), + "A job lock should exist at this point", ) def test_acquire_started_job(self): diff --git a/test_queue_job/tests/test_requeue_dead_job.py b/test_queue_job/tests/test_requeue_dead_job.py index a267c43c87..faec906bf0 100644 --- a/test_queue_job/tests/test_requeue_dead_job.py +++ b/test_queue_job/tests/test_requeue_dead_job.py @@ -1,6 +1,5 @@ # Copyright 2025 ACSONE SA/NV # License AGPL-3.0 or later (https://www.gnu.org/licenses/agpl). -from contextlib import closing from datetime import datetime, timedelta from odoo.tests import tagged @@ -13,35 +12,6 @@ @tagged("post_install", "-at_install") class TestRequeueDeadJob(JobCommonCase): - def get_locks(self, uuid, cr=None): - """ - Retrieve lock rows - """ - if cr is None: - cr = self.env.cr - - cr.execute( - """ - SELECT - queue_job_id - FROM - queue_job_lock - WHERE - queue_job_id IN ( - SELECT - id - FROM - queue_job - WHERE - uuid = %s - ) - FOR NO KEY UPDATE SKIP LOCKED - """, - [uuid], - ) - - return cr.fetchall() - def test_add_lock_record(self): queue_job = self._get_demo_job("test_started_job") self.assertEqual(len(queue_job), 1) @@ -50,22 +20,16 @@ def test_add_lock_record(self): job_obj.set_started() self.assertEqual(job_obj.state, "started") - locks = self.get_locks(job_obj.uuid) - - self.assertEqual(1, len(locks)) + self.assertFalse(self.is_job_locked(job_obj)) def test_lock(self): queue_job = self._get_demo_job("test_started_job") job_obj = Job.load(self.env, queue_job.uuid) job_obj.set_started() - job_obj.lock() - - with closing(self.env.registry.cursor()) as new_cr: - locks = self.get_locks(job_obj.uuid, new_cr) + job_obj.lock("started") - # Row should be locked - self.assertEqual(0, len(locks)) + self.assertTrue(self.is_job_locked(job_obj)) def test_requeue_dead_jobs(self): queue_job = self._get_demo_job("test_enqueued_job")