Procházet zdrojové kódy

Remove transfer schedules on transfer delete

Currently, the transfer schedules remained undeleted in the database
when a transfer is deleted. After deletion, getting the deleted
transfer's schedules would return a NotFound error, as the transfer was
deleted. This also means that the schedule transfers are not deletable
through the API either, so the database would keep having leaked rows in
the transfer schedules table.

Additionally, the transfer_cron service would still have cron jobs
configured for deleted transfers. Not only that, even on service
restart, cron jobs would still be created for deleted transfers, based
on the non-deleted transfer schedules.

Now, when we delete the transfer, we also delete its schedules.
Additionally, _get_transfer_schedules_filter now filsters out schedules
which have deleted transfers.
Claudiu Belu před 1 týdnem
rodič
revize
547c279baa

+ 18 - 0
coriolis/conductor/rpc/server.py

@@ -1474,8 +1474,26 @@ class ConductorServerEndpoint(object):
         transfer = self._get_transfer(ctxt, transfer_id)
         transfer = self._get_transfer(ctxt, transfer_id)
         self._check_transfer_running_executions(ctxt, transfer)
         self._check_transfer_running_executions(ctxt, transfer)
         self._check_delete_reservation_for_transfer(transfer)
         self._check_delete_reservation_for_transfer(transfer)
+        self._delete_transfer_schedules(ctxt, transfer_id)
         db_api.delete_transfer(ctxt, transfer_id)
         db_api.delete_transfer(ctxt, transfer_id)
 
 
+    def _delete_transfer_schedules(self, ctxt, transfer_id):
+        for schedule in db_api.get_transfer_schedules(ctxt, transfer_id=transfer_id):
+            try:
+                db_api.delete_transfer_schedule(
+                    ctxt,
+                    transfer_id,
+                    schedule.id,
+                    None,
+                    lambda ctxt, sched: self._cleanup_schedule_resources(ctxt, sched),
+                )
+            except exception.NotFound:
+                LOG.debug(
+                    "Schedule '%s' for Transfer '%s' was already deleted.",
+                    schedule.id,
+                    transfer_id,
+                )
+
     @transfer_synchronized
     @transfer_synchronized
     def delete_transfer_disks(self, ctxt, transfer_id):
     def delete_transfer_disks(self, ctxt, transfer_id):
         transfer = self._get_transfer(ctxt, transfer_id, include_task_info=True)
         transfer = self._get_transfer(ctxt, transfer_id, include_task_info=True)

+ 4 - 0
coriolis/db/api.py

@@ -112,6 +112,10 @@ def _get_transfer_schedules_filter(
             models.Transfer.project_id == context.project_id
             models.Transfer.project_id == context.project_id
         )
         )
 
 
+    # NOTE: a schedule whose Transfer has been deleted is an orphaned record. It should
+    # not be surfaced, it has no valid Transfer to run against.
+    sched_filter = sched_filter.filter(models.Transfer.deleted_at == null())
+
     if transfer_id:
     if transfer_id:
         sched_filter = sched_filter.filter(models.Transfer.id == transfer_id)
         sched_filter = sched_filter.filter(models.Transfer.id == transfer_id)
     if schedule_id:
     if schedule_id:

+ 51 - 0
coriolis/tests/conductor/rpc/test_server.py

@@ -1652,6 +1652,7 @@ class ConductorServerEndpointTestCase(test_base.CoriolisBaseTestCase):
         )
         )
 
 
     @mock.patch.object(db_api, 'delete_transfer')
     @mock.patch.object(db_api, 'delete_transfer')
+    @mock.patch.object(server.ConductorServerEndpoint, '_delete_transfer_schedules')
     @mock.patch.object(
     @mock.patch.object(
         server.ConductorServerEndpoint, '_check_delete_reservation_for_transfer'
         server.ConductorServerEndpoint, '_check_delete_reservation_for_transfer'
     )
     )
@@ -1664,6 +1665,7 @@ class ConductorServerEndpointTestCase(test_base.CoriolisBaseTestCase):
         mock_get_transfer,
         mock_get_transfer,
         mock_check_transfer_running_executions,
         mock_check_transfer_running_executions,
         mock_check_delete_reservation_for_transfer,
         mock_check_delete_reservation_for_transfer,
+        mock_delete_transfer_schedules,
         mock_delete_transfer,
         mock_delete_transfer,
     ):
     ):
         testutils.get_wrapped_function(self.server.delete_transfer)(
         testutils.get_wrapped_function(self.server.delete_transfer)(
@@ -1678,10 +1680,59 @@ class ConductorServerEndpointTestCase(test_base.CoriolisBaseTestCase):
         mock_check_delete_reservation_for_transfer.assert_called_once_with(
         mock_check_delete_reservation_for_transfer.assert_called_once_with(
             mock_get_transfer.return_value
             mock_get_transfer.return_value
         )
         )
+        mock_delete_transfer_schedules.assert_called_once_with(
+            mock.sentinel.context, mock.sentinel.transfer_id
+        )
         mock_delete_transfer.assert_called_once_with(
         mock_delete_transfer.assert_called_once_with(
             mock.sentinel.context, mock.sentinel.transfer_id
             mock.sentinel.context, mock.sentinel.transfer_id
         )
         )
 
 
+    @mock.patch.object(db_api, 'delete_transfer_schedule')
+    @mock.patch.object(db_api, 'get_transfer_schedules')
+    def test_delete_transfer_schedules(
+        self, mock_get_transfer_schedules, mock_delete_transfer_schedule
+    ):
+        schedule = mock.Mock(id=mock.sentinel.schedule_id)
+        mock_get_transfer_schedules.return_value = [schedule]
+
+        self.server._delete_transfer_schedules(
+            mock.sentinel.context, mock.sentinel.transfer_id
+        )
+
+        mock_get_transfer_schedules.assert_called_once_with(
+            mock.sentinel.context, transfer_id=mock.sentinel.transfer_id
+        )
+        mock_delete_transfer_schedule.assert_called_once_with(
+            mock.sentinel.context,
+            mock.sentinel.transfer_id,
+            mock.sentinel.schedule_id,
+            None,
+            mock.ANY,
+        )
+
+    @mock.patch.object(db_api, 'delete_transfer_schedule')
+    @mock.patch.object(db_api, 'get_transfer_schedules')
+    def test_delete_transfer_schedules_already_deleted(
+        self, mock_get_transfer_schedules, mock_delete_transfer_schedule
+    ):
+        schedule = mock.Mock(id=mock.sentinel.schedule_id)
+        mock_get_transfer_schedules.return_value = [schedule]
+        mock_delete_transfer_schedule.side_effect = exception.NotFound(
+            "No such schedule"
+        )
+
+        self.server._delete_transfer_schedules(
+            mock.sentinel.context, mock.sentinel.transfer_id
+        )
+
+        mock_delete_transfer_schedule.assert_called_once_with(
+            mock.sentinel.context,
+            mock.sentinel.transfer_id,
+            mock.sentinel.schedule_id,
+            None,
+            mock.ANY,
+        )
+
     @mock.patch.object(server.ConductorServerEndpoint, 'get_transfer_tasks_execution')
     @mock.patch.object(server.ConductorServerEndpoint, 'get_transfer_tasks_execution')
     @mock.patch.object(server.ConductorServerEndpoint, '_begin_tasks')
     @mock.patch.object(server.ConductorServerEndpoint, '_begin_tasks')
     @mock.patch.object(db_api, "add_transfer_tasks_execution")
     @mock.patch.object(db_api, "add_transfer_tasks_execution")

+ 28 - 0
coriolis/tests/db/test_api.py

@@ -936,6 +936,34 @@ class TransferSchedulesDBAPITestCase(BaseDBAPITestCase):
         ).first()
         ).first()
         self.assertEqual(result, self.valid_transfer_schedule)
         self.assertEqual(result, self.valid_transfer_schedule)
 
 
+    def test__get_transfer_schedules_filter_excludes_deleted_transfer(self):
+        deleted_transfer = models.Transfer()
+        deleted_transfer.id = str(uuid.uuid4())
+        deleted_transfer.user_id = "1"
+        deleted_transfer.project_id = "1"
+        deleted_transfer.base_id = deleted_transfer.id
+        deleted_transfer.scenario = constants.TRANSFER_SCENARIO_REPLICA
+        deleted_transfer.last_execution_status = DEFAULT_EXECUTION_STATUS
+        deleted_transfer.executions = []
+        deleted_transfer.instances = [DEFAULT_INSTANCE]
+        deleted_transfer.info = DEFAULT_TASK_INFO
+        deleted_transfer.origin_endpoint_id = self.valid_transfer.origin_endpoint_id
+        deleted_transfer.destination_endpoint_id = (
+            self.valid_transfer.destination_endpoint_id
+        )
+        deleted_transfer.deleted_at = timeutils.utcnow()
+        self.session.add(deleted_transfer)
+
+        orphaned_schedule = self._create_dummy_transfer_schedule(
+            deleted_transfer, expiration_date=None
+        )
+        self.session.add(orphaned_schedule)
+
+        result = api._get_transfer_schedules_filter(
+            self.context, schedule_id=orphaned_schedule.id
+        ).first()
+        self.assertIsNone(result)
+
     def test__get_transfer_schedules_filter_by_not_expired(self):
     def test__get_transfer_schedules_filter_by_not_expired(self):
         expiration_date = timeutils.utcnow() + datetime.timedelta(days=1)
         expiration_date = timeutils.utcnow() + datetime.timedelta(days=1)
         unexpired_transfer_schedule = self._create_dummy_transfer_schedule(
         unexpired_transfer_schedule = self._create_dummy_transfer_schedule(

+ 1 - 1
coriolis/tests/integration/base.py

@@ -160,7 +160,7 @@ class CoriolisIntegrationTestBase(test_base.CoriolisBaseTestCase):
             skip_os_morphing=True,
             skip_os_morphing=True,
             **kwargs,
             **kwargs,
         )
         )
-        self.addCleanup(self._client.transfers.delete, transfer.id)
+        self.addCleanup(self._ignoreExc(self._client.transfers.delete), transfer.id)
 
 
         return transfer
         return transfer
 
 

+ 22 - 0
coriolis/tests/integration/transfers/test_schedules.py

@@ -12,6 +12,7 @@ import time
 
 
 from oslo_utils import timeutils
 from oslo_utils import timeutils
 
 
+from coriolis.db import api as db_api
 from coriolis.tests.integration import base
 from coriolis.tests.integration import base
 
 
 
 
@@ -137,3 +138,24 @@ class TransferScheduleTests(_TransferScheduleTestBase):
         )
         )
 
 
         self.assertExecutionCompleted(execution.id)
         self.assertExecutionCompleted(execution.id)
+
+    def test_deleting_transfer_deletes_schedule(self):
+        """Deleting a Transfer must delete its schedules along with it.
+
+        Regression test: ``delete_transfer`` used to leave attached schedules in place,
+        both in the DB and registered with the transfer-cron service, which kept firing
+        on its configured cadence against a Transfer that no longer existed.
+        """
+        target = timeutils.utcnow() + datetime.timedelta(seconds=10)
+        self._create_schedule(
+            schedule={"minute": target.minute, "hour": target.hour},
+            enabled=True,
+        )
+
+        self._client.transfers.delete(self._transfer.id)
+
+        ctxt = self._get_db_context()
+        remaining = db_api.get_transfer_schedules(ctxt, transfer_id=self._transfer.id)
+        self.assertEqual(
+            [], remaining, "Schedule was not deleted along with its Transfer"
+        )