Skip to content
Merged
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
37 changes: 35 additions & 2 deletions coriolis/tests/transfer_cron/rpc/test_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
{
Expand Down
47 changes: 43 additions & 4 deletions coriolis/transfer_cron/rpc/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
import json
import os
import threading
import time

from oslo_log import log as logging
from oslo_utils import timeutils
Expand Down Expand Up @@ -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")
Expand Down Expand Up @@ -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(
Expand Down
Loading