from collections.abc import AsyncIterator import asyncio from cloudcoil.client import Config from cloudcoil.errors import ResourceNotFound from dbus_fast.aio import MessageBus from grpclib.client import Channel import cloudcoil.models.kubernetes as k8s import cloudcoil.models.volumesnapshot as k8s_vs import lvmd.proto as lvmd import pytest from .testutils import VM_SOCK_DIR from .testutils.k8s import KUBECONFIG, TEST_NAMESPACE from .testutils.stratis import Stratisd from .testutils.systemd import Systemd _OPTION_SKIP_CLEANUP = '--skip-cleanup' def pytest_addoption(parser: pytest.Parser) -> None: parser.addoption( _OPTION_SKIP_CLEANUP, action='store_true', default=False, help='Skip cleaning up the k8s cluster / stratis pools', ) @pytest.fixture async def bus() -> AsyncIterator[MessageBus]: bus = MessageBus(bus_address=f'unix:path={VM_SOCK_DIR}/dbus.sock') _ = await bus.connect() yield bus bus.disconnect() async with asyncio.timeout(5): await bus.wait_for_disconnect() @pytest.fixture async def stratisd(bus: MessageBus) -> Stratisd: return await Stratisd.create(bus) @pytest.fixture async def systemd(bus: MessageBus) -> Systemd: return await Systemd.create(bus) async def _cleanup(stratisd: Stratisd) -> None: async with asyncio.timeout(60): try: res = await k8s.core.v1.Namespace.async_delete(TEST_NAMESPACE) if isinstance(res, k8s.core.v1.Namespace): await res.async_wait_for(lambda event, _: event == 'DELETED') except ResourceNotFound: pass for kind in ( k8s.core.v1.PersistentVolume, k8s_vs.snapshot.v1.VolumeSnapshotContent, k8s_vs.snapshot.v1.VolumeSnapshotContent, k8s_vs.snapshot.v1.VolumeSnapshotClass, ): for resource in await kind.async_delete_all(): await resource.async_wait_for(lambda event, _: event == 'DELETED') await stratisd.reload() await stratisd.wipe_all_filesystems() @pytest.fixture(autouse=True) async def k8s_config( stratisd: Stratisd, request: pytest.FixtureRequest ) -> AsyncIterator[Config]: with Config(kubeconfig=KUBECONFIG, namespace=TEST_NAMESPACE) as config: await _cleanup(stratisd) yield config if not request.config.getoption(_OPTION_SKIP_CLEANUP): await _cleanup(stratisd) @pytest.fixture async def lvmd_channel() -> Channel: return Channel(path=str(VM_SOCK_DIR / 'lvmd.sock')) @pytest.fixture def lvservice(lvmd_channel: Channel) -> lvmd.LVServiceStub: return lvmd.LVServiceStub(lvmd_channel) @pytest.fixture def vgservice(lvmd_channel: Channel) -> lvmd.VGServiceStub: return lvmd.VGServiceStub(lvmd_channel)