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
24 changes: 9 additions & 15 deletions queue_job/controllers/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -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__)

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

Expand Down Expand Up @@ -61,29 +61,23 @@ 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",
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()
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,
Expand All @@ -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
Expand Down
66 changes: 17 additions & 49 deletions queue_job/job.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand All @@ -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):
Expand Down Expand Up @@ -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
Expand Down
89 changes: 26 additions & 63 deletions queue_job/jobrunner/runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
"""
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
10 changes: 10 additions & 0 deletions test_queue_job/tests/common.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
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
42 changes: 3 additions & 39 deletions test_queue_job/tests/test_requeue_dead_job.py
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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)
Expand All @@ -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")
Expand Down
Loading