diff --git a/coriolis/tests/transfer_cron/rpc/test_server.py b/coriolis/tests/transfer_cron/rpc/test_server.py index 15c2fce47..d89b08926 100644 --- a/coriolis/tests/transfer_cron/rpc/test_server.py +++ b/coriolis/tests/transfer_cron/rpc/test_server.py @@ -90,10 +90,43 @@ def test_init_defers_cron_to_after_fork( @mock.patch.object(server.TransferCronServerEndpoint, '_init_cron') def test_ensure_cron_started_is_idempotent(self, mock_init_cron): self.server._cron_started = False - self.server._ensure_cron_started() - self.server._ensure_cron_started() + self.server._loader_started = False + + def _run_inline(target, args=(), kwargs=None, daemon=True): + target(*(args or ()), **(kwargs or {})) + + with mock.patch.object(server.utils, 'start_thread', side_effect=_run_inline): + self.server._ensure_cron_started() + self.server._ensure_cron_started() mock_init_cron.assert_called_once() + self.assertTrue(self.server._cron_started) + + @mock.patch.object(server.time, 'sleep') + def test_load_schedules_retries_after_conductor_timeout(self, mock_sleep): + self.server._cron_started = False + self.server._init_cron = mock.Mock(side_effect=[Exception('timeout'), None]) + + with self.assertLogs('coriolis.transfer_cron.rpc.server', level=logging.ERROR): + self.server._load_schedules_until_started() + + self.assertEqual(self.server._init_cron.call_count, 2) + mock_sleep.assert_called_once_with(60) + self.assertTrue(self.server._cron_started) + + def test_ensure_cron_started_does_not_raise(self): + self.server._cron_started = False + self.server._loader_started = False + + with mock.patch.object( + server.utils, 'start_thread', side_effect=RuntimeError('boom') + ): + with self.assertLogs( + 'coriolis.transfer_cron.rpc.server', level=logging.ERROR + ): + self.server._ensure_cron_started() + + self.assertFalse(self.server._loader_started) @ddt.data( { diff --git a/coriolis/transfer_cron/rpc/server.py b/coriolis/transfer_cron/rpc/server.py index c02383816..bc0ef886c 100644 --- a/coriolis/transfer_cron/rpc/server.py +++ b/coriolis/transfer_cron/rpc/server.py @@ -4,6 +4,7 @@ import json import os import threading +import time from oslo_log import log as logging from oslo_utils import timeutils @@ -44,19 +45,57 @@ def __init__(self): self._admin_ctx = context.get_admin_context() self._cron_lock = threading.Lock() self._cron_started = False + self._cron_loop_started = False + self._loader_started = False # NOTE (fabi200123): oslo_service forks worker processes even when # workers=1. The cron loop must run in the same process as the RPC # handlers that register/unregister jobs, otherwise the loop checks # a job registry in the parent process while registrations land in # the forked child. Defer cron startup until after the fork. + # + # The callback itself must not raise. CPython ignores exceptions + # from after_in_child ("Exception ignored in: ..."), so a + # MessagingTimeout while the conductor is still starting used to + # leave this process up with the schedule loop never started. os.register_at_fork(after_in_child=self._ensure_cron_started) def _ensure_cron_started(self): + try: + with self._cron_lock: + if self._cron_started or self._loader_started: + return + self._loader_started = True + self._loader_thread = utils.start_thread(self._load_schedules_until_started) + except Exception: + LOG.exception("Failed to start the transfer schedule loader") + with self._cron_lock: + if not self._cron_started: + self._loader_started = False + + def _load_schedules_until_started(self): + delay = 60 + while True: + try: + self._init_cron() + except Exception: + LOG.exception( + "Failed to load transfer schedules from the conductor. " + "Retrying in %s seconds", + delay, + ) + time.sleep(delay) + continue + with self._cron_lock: + self._cron_started = True + LOG.info("Transfer schedules loaded, cron loop started") + return + + def _start_cron_loop(self): with self._cron_lock: - if self._cron_started: + if self._cron_loop_started: return - self._init_cron() - self._cron_started = True + self._cron.start() + self._cron_loop_started = True def _deserialize_schedule(self, sched): expires = sched.get("expiration_date") @@ -107,7 +146,7 @@ def _init_cron(self): # of an invalid schedule that managed to creep its # way into the DB, or just ignore that one schedule? LOG.exception(err) - self._cron.start() + self._start_cron_loop() def _get_all_schedules(self): schedules = self._rpc_client.get_transfer_schedules(