diff --git a/queue_job/job.py b/queue_job/job.py index a69a71ddc5..1251ca5478 100644 --- a/queue_job/job.py +++ b/queue_job/job.py @@ -59,6 +59,9 @@ def _rebind_to_cr(value, cr): DEFAULT_MAX_RETRIES = 5 RETRY_INTERVAL = 10 * 60 # seconds +# struct.unpack('i', hashlib.sha1(b"odoo:queue.job").digest()[:4]) +QUEUE_JOB_LOCK_KEY = 2129646392 # int4: -2**31 <= ... <= 2**31 - 1 + _logger = logging.getLogger(__name__) @@ -243,55 +246,25 @@ 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. + """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. + 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" + # Even if 'id' is INTEGER type, be safe and apply a modulo + " AND pg_try_advisory_xact_lock(%s, id %% (1<<31));" ) - - # 1 job should be locked - return bool(self.env.cr.fetchall()) + self.env.cr.execute(lock_query, [self.uuid, STARTED, QUEUE_JOB_LOCK_KEY]) + if not self.env.cr.fetchone(): + _logger.debug("Lock NOT acquired on Job %s", self.uuid) + return False + _logger.debug("Lock acquired on Job %s", self.uuid) + return True @classmethod def _load_from_db_record(cls, job_db_record): @@ -852,7 +825,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..65e7def52e 100644 --- a/queue_job/jobrunner/runner.py +++ b/queue_job/jobrunner/runner.py @@ -35,6 +35,7 @@ import odoo from odoo.tools import config +from ..job import QUEUE_JOB_LOCK_KEY from . import queue_job_config from .channels import ENQUEUED, NOT_DONE, RELOAD_PAYLOAD, ChannelConfig, ChannelManager @@ -354,69 +355,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' + AND pg_try_advisory_xact_lock(%s, id %% (1<<31))) + + 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): """ @@ -447,7 +411,7 @@ def requeue_dead_jobs(self): with closing(self.conn.cursor()) as cr: query = self._query_requeue_dead_jobs() - cr.execute(query) + cr.execute(query, (QUEUE_JOB_LOCK_KEY,)) for (uuid,) in cr.fetchall(): _logger.warning("Re-queued dead job with uuid: %s", uuid) 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..eda7dbf0c8 100644 --- a/test_queue_job/tests/common.py +++ b/test_queue_job/tests/common.py @@ -3,7 +3,7 @@ from odoo.tests import common -from odoo.addons.queue_job.job import Job +from odoo.addons.queue_job.job import QUEUE_JOB_LOCK_KEY, Job class JobCommonCase(common.TransactionCase): @@ -30,3 +30,17 @@ def _get_demo_job(self, uuid): "to make this test work", ) return job + + def is_job_locked(self, job, cr=None): + base_query = "SELECT 1 FROM queue_job WHERE uuid = %s" + query_lock_shared, query_unlock_shared = ( + f"{base_query} AND pg_try_advisory_lock_shared(%s, id);", + f"{base_query} AND pg_advisory_unlock_shared(%s, id);", + ) + args = job.uuid, QUEUE_JOB_LOCK_KEY + with self.env.registry.cursor() as cr: + cr.execute(query_lock_shared, args) + if not cr.fetchone(): + return True + cr.execute(query_unlock_shared, args) + 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..cea0768ebd 100644 --- a/test_queue_job/tests/test_requeue_dead_job.py +++ b/test_queue_job/tests/test_requeue_dead_job.py @@ -1,11 +1,10 @@ # 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 -from odoo.addons.queue_job.job import Job +from odoo.addons.queue_job.job import QUEUE_JOB_LOCK_KEY, Job from odoo.addons.queue_job.jobrunner.runner import Database from .common import JobCommonCase @@ -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,9 +20,7 @@ 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") @@ -61,11 +29,7 @@ def test_lock(self): 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) - - # 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") @@ -78,7 +42,7 @@ def test_requeue_dead_jobs(self): # requeue dead jobs using current cursor query = Database(self.env.cr.dbname)._query_requeue_dead_jobs() - self.env.cr.execute(query) + self.env.cr.execute(query, (QUEUE_JOB_LOCK_KEY,)) uuids_requeued = self.env.cr.fetchall() self.assertTrue(queue_job.uuid in j[0] for j in uuids_requeued) @@ -95,6 +59,6 @@ def test_requeue_orphaned_jobs(self): # job is now picked up by the requeue query (which includes orphaned jobs) query = Database(self.env.cr.dbname)._query_requeue_dead_jobs() - self.env.cr.execute(query) + self.env.cr.execute(query, (QUEUE_JOB_LOCK_KEY,)) uuids_requeued = self.env.cr.fetchall() self.assertTrue(queue_job.uuid in j[0] for j in uuids_requeued)