From 2755af8fac13d2fd5ddc7d90ad3a633ab4c8bd2f Mon Sep 17 00:00:00 2001 From: Julien Castiaux Date: Wed, 27 Jul 2022 14:49:29 +0000 Subject: [PATCH] [IMP] base: test cases for ir_cron The cron subsystem is the system responsible of running background task at regular interval, it runs in multiple dedicated threads or workers that are independent of the regular HTTP threads/workers. There was a major overhaul of the system in v15 (4b28f1162a8) to introduce cron triggers, a way to run a task at a given moment in addition to the regular configured interval. Although the cron system was not extensively tested before that v15 refactor, no new test were introduced with that refactor leaving the system mostly untested. Since then we had to fix multiple subtle concurrency bugs such as b940d1c25f8 and c06cee44fe1. Due to the lack of an existing test suite, no regression tests were added next to those fixes. With this commit we introduce the missing cron test suite. The test suite is separated in two different test cases: - A standard pre-install TransactionCase to test everything that can be tested with a single cursor. This case can be run the usual way with `--test-tags :TestIrCron`. - A non-standard post-install **database breaking** test case to run concurrency tests that often require multiple SQL transactions. This case requires the special `--test-tags database_breaking` to be executed. You MUST backup your current database before running that test or you'll loose data. closes odoo/odoo#97087 Signed-off-by: Julien Castiaux --- odoo/addons/base/tests/test_ir_cron.py | 345 ++++++++++++++++++++++--- 1 file changed, 315 insertions(+), 30 deletions(-) diff --git a/odoo/addons/base/tests/test_ir_cron.py b/odoo/addons/base/tests/test_ir_cron.py index 68c83ed97d0..d2e7fef056f 100644 --- a/odoo/addons/base/tests/test_ir_cron.py +++ b/odoo/addons/base/tests/test_ir_cron.py @@ -1,10 +1,19 @@ # -*- coding: utf-8 -*- # Part of Odoo. See LICENSE file for full copyright and licensing details. -from unittest.mock import patch +import collections +import secrets +import textwrap +import threading +from concurrent.futures import ThreadPoolExecutor +from datetime import timedelta +from unittest.mock import call, patch +from freezegun import freeze_time -from odoo import fields -from odoo.tests.common import TransactionCase, RecordCapturer +import odoo +from odoo import api, fields +from odoo.tests.common import BaseCase, TransactionCase, RecordCapturer, get_db_name, tagged +from odoo.tools import mute_logger class CronMixinCase: @@ -29,40 +38,316 @@ class CronMixinCase: domain=[('cron_id', '=', cron_id)] if cron_id else [] ) - -class TestIrCron(TransactionCase, CronMixinCase): - - def setUp(self): - super(TestIrCron, self).setUp() - - self.cron = self.env['ir.cron'].create({ - 'name': 'TestCron', - 'model_id': self.env.ref('base.model_res_partner').id, + @classmethod + def _get_cron_data(cls, env, priority=5): + unique = secrets.token_urlsafe(8) + return { + 'name': f'Dummy cron for TestIrCron {unique}', 'state': 'code', - 'code': 'model.search([("name", "=", "TestCronRecord")]).write({"name": "You have been CRONWNED"})', + 'code': '', + 'model_id': env.ref('base.model_res_partner').id, + 'model_name': 'res.partner', + 'user_id': env.uid, + 'active': True, 'interval_number': 1, 'interval_type': 'days', 'numbercall': -1, 'doall': False, - }) - self.test_partner = self.env['res.partner'].create({ - 'name': 'TestCronRecord' - }) - self.test_partner2 = self.env['res.partner'].create({ - 'name': 'NotTestCronRecord' - }) + 'nextcall': fields.Datetime.now() + timedelta(hours=1), + 'lastcall': False, + 'priority': priority, + } + + @classmethod + def _get_partner_data(cls, env): + unique = secrets.token_urlsafe(8) + return {'name': f'Dummy partner for TestIrCron {unique}'} + + +class TestIrCron(TransactionCase, CronMixinCase): + + @classmethod + def setUpClass(cls): + super().setUpClass() + + freezer = freeze_time(cls.cr.now()) + freezer.start() + cls.addClassCleanup(freezer.stop) + + cls.cron = cls.env['ir.cron'].create(cls._get_cron_data(cls.env)) + cls.partner = cls.env['res.partner'].create(cls._get_partner_data(cls.env)) + + def setUp(self): + self.partner.write(self._get_partner_data(self.env)) + self.cron.write(self._get_cron_data(self.env)) + self.env['ir.cron.trigger'].search( + [('cron_id', '=', self.cron.id)] + ).unlink() def test_cron_direct_trigger(self): - self.assertFalse(self.cron.lastcall) - self.assertEqual(self.test_partner.name, 'TestCronRecord') - self.assertEqual(self.test_partner2.name, 'NotTestCronRecord') + self.cron.code = textwrap.dedent(f"""\ + model.search( + [("id", "=", {self.partner.id})] + ).write( + {{"name": "You have been CRONWNED"}} + ) + """) - def patched_now(*args, **kwargs): - return '2020-10-22 08:00:00' + self.cron.method_direct_trigger() - with patch('odoo.fields.Datetime.now', patched_now): - self.cron.method_direct_trigger() + self.assertEqual(self.cron.lastcall, fields.Datetime.now()) + self.assertEqual(self.partner.name, 'You have been CRONWNED') - self.assertEqual(fields.Datetime.to_string(self.cron.lastcall), '2020-10-22 08:00:00') - self.assertEqual(self.test_partner.name, 'You have been CRONWNED') - self.assertEqual(self.test_partner2.name, 'NotTestCronRecord') + def test_cron_no_job_ready(self): + self.cron.nextcall = fields.Datetime.now() + timedelta(days=1) + self.cron.flush_recordset() + + ready_jobs = self.registry['ir.cron']._get_all_ready_jobs(self.cr) + self.assertNotIn(self.cron.id, [job['id'] for job in ready_jobs]) + + def test_cron_ready_by_nextcall(self): + self.cron.nextcall = fields.Datetime.now() + self.cron.flush_recordset() + + ready_jobs = self.registry['ir.cron']._get_all_ready_jobs(self.cr) + self.assertIn(self.cron.id, [job['id'] for job in ready_jobs]) + + def test_cron_ready_by_trigger(self): + self.cron._trigger() + self.env['ir.cron.trigger'].flush_model() + + ready_jobs = self.registry['ir.cron']._get_all_ready_jobs(self.cr) + self.assertIn(self.cron.id, [job['id'] for job in ready_jobs]) + + def test_cron_unactive_never_ready(self): + self.cron.active = False + self.cron.nextcall = fields.Datetime.now() + self.cron._trigger() + self.cron.flush_recordset() + self.env['ir.cron.trigger'].flush_model() + + ready_jobs = self.registry['ir.cron']._get_all_ready_jobs(self.cr) + self.assertNotIn(self.cron.id, [job['id'] for job in ready_jobs]) + + def test_cron_numbercall0_never_ready(self): + self.cron.numbercall = 0 + self.cron.nextcall = fields.Datetime.now() + self.cron._trigger() + self.cron.flush_recordset() + self.env['ir.cron.trigger'].flush_model() + + ready_jobs = self.registry['ir.cron']._get_all_ready_jobs(self.cr) + self.assertNotIn(self.cron.id, [job['id'] for job in ready_jobs]) + + def test_cron_ready_jobs_order(self): + cron_avg = self.cron.copy() + cron_avg.priority = 5 # average priority + + cron_high = self.cron.copy() + cron_high.priority = 0 # highest priority + + cron_low = self.cron.copy() + cron_low.priority = 10 # lowest priority + + crons = cron_high | cron_avg | cron_low # order is important + crons.write({'nextcall': fields.Datetime.now()}) + crons.flush_recordset() + ready_jobs = self.registry['ir.cron']._get_all_ready_jobs(self.cr) + + self.assertEqual( + [job['id'] for job in ready_jobs if job['id'] in crons._ids], + list(crons._ids), + ) + + def test_cron_process_job(self): + + Setup = collections.namedtuple('Setup', ['doall', 'numbercall', 'missedcall', 'trigger']) + Expect = collections.namedtuple('Expect', ['call_count', 'call_left', 'active']) + + matrix = [ + (Setup(doall=False, numbercall=-1, missedcall=2, trigger=False), + Expect(call_count=1, call_left=-1, active=True)), + (Setup(doall=True, numbercall=-1, missedcall=2, trigger=False), + Expect(call_count=2, call_left=-1, active=True)), + (Setup(doall=False, numbercall=3, missedcall=2, trigger=False), + Expect(call_count=1, call_left=2, active=True)), + (Setup(doall=True, numbercall=3, missedcall=2, trigger=False), + Expect(call_count=2, call_left=1, active=True)), + (Setup(doall=True, numbercall=3, missedcall=4, trigger=False), + Expect(call_count=3, call_left=0, active=False)), + (Setup(doall=True, numbercall=3, missedcall=0, trigger=True), + Expect(call_count=1, call_left=2, active=True)), + ] + + for setup, expect in matrix: + with self.subTest(setup=setup, expect=expect): + self.cron.write({ + 'active': True, + 'doall': setup.doall, + 'numbercall': setup.numbercall, + 'nextcall': fields.Datetime.now() - timedelta(days=setup.missedcall - 1), + }) + with self.capture_triggers(self.cron.id) as capture: + if setup.trigger: + self.cron._trigger() + + self.cron.flush_recordset() + capture.records.flush_recordset() + self.registry.enter_test_mode(self.cr) + try: + with patch.object(self.registry['ir.cron'], '_callback') as callback: + self.registry['ir.cron']._process_job( + self.registry.db_name, + self.registry.cursor(), + self.cron.read(load=None)[0] + ) + finally: + self.registry.leave_test_mode() + self.cron.invalidate_recordset() + capture.records.invalidate_recordset() + + self.assertEqual(callback.call_count, expect.call_count) + self.assertEqual(self.cron.numbercall, expect.call_left) + self.assertEqual(self.cron.active, expect.active) + self.assertEqual(self.cron.lastcall, fields.Datetime.now()) + self.assertEqual(self.cron.nextcall, fields.Datetime.now() + timedelta(days=1)) + self.assertEqual(self.env['ir.cron.trigger'].search_count([ + ('cron_id', '=', self.cron.id), + ('call_at', '<=', fields.Datetime.now())] + ), 0) + + +@tagged('-standard', '-at_install', 'post_install', 'database_breaking') +class TestIrCronConcurrent(BaseCase, CronMixinCase): + + @classmethod + def setUpClass(cls): + super().setUpClass() + + # Keep a reference on the real cron methods, those without patch + cls.registry = odoo.registry(get_db_name()) + cls.cron_process_job = cls.registry['ir.cron']._process_job + cls.cron_process_jobs = cls.registry['ir.cron']._process_jobs + cls.cron_get_all_ready_jobs = cls.registry['ir.cron']._get_all_ready_jobs + cls.cron_acquire_one_job = cls.registry['ir.cron']._acquire_one_job + cls.cron_callback = cls.registry['ir.cron']._callback + + def setUp(self): + super().setUp() + + with self.registry.cursor() as cr: + env = api.Environment(cr, odoo.SUPERUSER_ID, {}) + env['ir.cron'].search([]).unlink() + env['ir.cron.trigger'].search([]).unlink() + + self.cron1_data = env['ir.cron'].create(self._get_cron_data(env, priority=1)).read(load=None)[0] + self.cron2_data = env['ir.cron'].create(self._get_cron_data(env, priority=2)).read(load=None)[0] + self.partner_data = env['res.partner'].create(self._get_partner_data(env)).read(load=None)[0] + self.cron_ids = [self.cron1_data['id'], self.cron2_data['id']] + + def test_cron_concurrency_1(self): + """ + Two cron threads "th1" and "th2" wake up at the same time and + see two jobs "job1" and "job2" that are ready (setup). + + Th1 acquire job1, before it can process and release its job, th2 + acquire a job too (setup). Th2 shouldn't be able to acquire job1 + as another thread is processing it, it should skips job1 and + should acquire job2 instead (test). Both thread then process + their job, update its `nextcall` and release it (setup). + + All the threads update and release their job before any thread + attempt to acquire another job. (setup) + + The two thread each attempt to acquire a new job (setup), they + should both fail to acquire any as each job's nextcall is in the + future* (test). + + *actually, in their own transaction, the other job's nextcall is + still "in the past" but any attempt to use that information + would result in a serialization error. This tests ensure that + that serialization error is correctly handled and ignored. + """ + lock = threading.Lock() + barrier = threading.Barrier(2) + + ### + # Setup + ### + + # Watchdog, if a thread was waiting at the barrier when the + # other exited, it receives a BrokenBarrierError and exits too. + def process_jobs(*args, **kwargs): + try: + self.cron_process_jobs(*args, **kwargs) + finally: + barrier.reset() + + # The two threads get the same list of jobs + def get_all_ready_jobs(*args, **kwargs): + jobs = self.cron_get_all_ready_jobs(*args, **kwargs) + barrier.wait() + return jobs + + # When a thread acquire a job, it processes it till the end + # before another thread can acquire one. + def acquire_one_job(*args, **kwargs): + lock.acquire(timeout=1) + try: + with mute_logger('odoo.sql_db'): + job = self.cron_acquire_one_job(*args, **kwargs) + except Exception: + lock.release() + raise + if not job: + lock.release() + return job + + # When a thread is done processing its job, it waits for the + # other thread to catch up. + def process_job(*args, **kwargs): + try: + return_value = self.cron_process_job(*args, **kwargs) + finally: + lock.release() + barrier.wait(timeout=1) + return return_value + + # Set 2 jobs ready, process them in 2 different threads. + with self.registry.cursor() as cr: + env = api.Environment(cr, odoo.SUPERUSER_ID, {}) + env['ir.cron'].browse(self.cron_ids).write({ + 'nextcall': fields.Datetime.now() - timedelta(hours=1) + }) + + ### + # Run + ### + with patch.object(self.registry['ir.cron'], '_process_jobs', process_jobs), \ + patch.object(self.registry['ir.cron'], '_get_all_ready_jobs', get_all_ready_jobs), \ + patch.object(self.registry['ir.cron'], '_acquire_one_job', acquire_one_job), \ + patch.object(self.registry['ir.cron'], '_process_job', process_job), \ + patch.object(self.registry['ir.cron'], '_callback') as callback, \ + ThreadPoolExecutor(max_workers=2) as executor: + fut1 = executor.submit(self.registry['ir.cron']._process_jobs, self.registry.db_name) + fut2 = executor.submit(self.registry['ir.cron']._process_jobs, self.registry.db_name) + fut1.result(timeout=2) + fut2.result(timeout=2) + + ### + # Validation + ### + + self.assertEqual(len(callback.call_args_list), 2, 'Two jobs must have been processed.') + self.assertEqual(callback.call_args_list, [ + call( + self.cron1_data['name'], + self.cron1_data['ir_actions_server_id'], + self.cron1_data['id'], + ), + call( + self.cron2_data['name'], + self.cron2_data['ir_actions_server_id'], + self.cron2_data['id'], + ), + ])