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
58 changes: 14 additions & 44 deletions queue_job/job.py
Original file line number Diff line number Diff line change
Expand Up @@ -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**32 <= ... <= 2**32 - 1

_logger = logging.getLogger(__name__)


Expand Down Expand Up @@ -243,55 +246,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.
"""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 AND pg_try_advisory_xact_lock(%s, id);"
)

# 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):
Expand Down Expand Up @@ -852,7 +823,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
Expand Down
92 changes: 28 additions & 64 deletions queue_job/jobrunner/runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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))

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):
"""
Expand Down Expand Up @@ -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)
Expand Down
1 change: 0 additions & 1 deletion queue_job/models/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,4 +3,3 @@
from . import queue_job
from . import queue_job_channel
from . import queue_job_function
from . import queue_job_lock
16 changes: 0 additions & 16 deletions queue_job/models/queue_job_lock.py

This file was deleted.

1 change: 0 additions & 1 deletion queue_job/security/ir.model.access.csv
Original file line number Diff line number Diff line change
@@ -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
Expand Down
16 changes: 15 additions & 1 deletion test_queue_job/tests/common.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand All @@ -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
12 changes: 4 additions & 8 deletions test_queue_job/tests/test_acquire_job.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand 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):
Expand Down
46 changes: 5 additions & 41 deletions test_queue_job/tests/test_requeue_dead_job.py
Original file line number Diff line number Diff line change
@@ -1,47 +1,17 @@
# 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


@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)
Expand All @@ -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")
Expand All @@ -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")
Expand All @@ -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)
Expand All @@ -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)
Loading