diff --git a/pyproject.toml b/pyproject.toml
index 285df60..3dd9357 100644
--- a/pyproject.toml
+++ b/pyproject.toml
@@ -24,7 +24,10 @@ dev = [
[tool.mypy]
python_version = "3.14"
-mypy_path = "stubs"
+mypy_path = [
+ "src",
+ "stubs"
+]
files = [
"src",
"tests"
@@ -32,6 +35,7 @@ files = [
disallow_untyped_defs = true
enable_error_code = ["explicit-override"]
+explicit_package_bases = true
ignore_missing_imports = false
implicit_optional = false
diff --git a/src/xcp_storage/backends/__init__.py b/src/xcp_storage/backends/__init__.py
new file mode 100644
index 0000000..e69de29
diff --git a/src/xcp_storage/backends/drbd/__init__.py b/src/xcp_storage/backends/drbd/__init__.py
new file mode 100644
index 0000000..ede438f
--- /dev/null
+++ b/src/xcp_storage/backends/drbd/__init__.py
@@ -0,0 +1,234 @@
+# Copyright (C) 2026 Vates SAS
+#
+# This program is free software: you can redistribute it and/or modify
+# it under the terms of the GNU General Public License as published by
+# the Free Software Foundation, either version 3 of the License, or
+# (at your option) any later version.
+# This program is distributed in the hope that it will be useful,
+# but WITHOUT ANY WARRANTY; without even the implied warranty of
+# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
+# GNU General Public License for more details.
+#
+# You should have received a copy of the GNU General Public License
+# along with this program. If not, see .
+
+import contextlib
+from dataclasses import dataclass
+import json
+from pathlib import Path
+import re
+
+import xcp_storage.log as log
+from xcp_storage.utils.process import (
+ get_process_cmdline,
+ run_command,
+)
+
+from xcp_storage.typing import (
+ Any,
+ Dict,
+ Final,
+ Iterator,
+ List,
+)
+
+# ==============================================================================
+
+logger = log.get_logger() # Use default logger.
+
+# ------------------------------------------------------------------------------
+
+DRBD_BY_RES_PATH: Final = "/dev/drbd/by-res/"
+
+DRBD_PORT_RANGE: Final = (7000, 8000)
+
+# ------------------------------------------------------------------------------
+
+_EXEC_PATH_DRBDSETUP: Final = "/usr/sbin/drbdsetup"
+
+_REGEX_DRBD_OPENER_LINE: Final = re.compile(r"(.*)\s+(\d+)\s+(\d+)")
+
+# Characters allowed in a resource name: ASCII only, and neither a path separator nor a leading `-` or `.`.
+_REGEX_DRBD_RESOURCE_NAME: Final = re.compile(r"\w[\w.-]*", re.ASCII)
+
+# ------------------------------------------------------------------------------
+
+# Because this module can be used by external layers and RPC, we must have checkers
+# to prevent code injection.
+
+def _check_drbd_resource_name(resource_name: str) -> None:
+ if not isinstance(resource_name, str):
+ raise ValueError(f"Not a DRBD resource name: `{resource_name}`.")
+ if not _REGEX_DRBD_RESOURCE_NAME.fullmatch(resource_name):
+ raise ValueError(f"Invalid DRBD resource name: `{resource_name}`.")
+
+def _check_drbd_volume_number(volume_number: int) -> None:
+ if not isinstance(volume_number, int) or isinstance(volume_number, bool):
+ raise ValueError(f"Not a DRBD volume number: `{volume_number}`.")
+ if volume_number < 0:
+ raise ValueError(f"Invalid DRBD volume number: `{volume_number}`.")
+
+# ------------------------------------------------------------------------------
+
+@contextlib.contextmanager
+def _handle_drbd_json_error() -> Iterator[None]:
+ """
+ Log instead of raising an error while reading the JSON status of DRBD.
+
+ The log message suggests that the JSON format may have changed BUT it's also logged for valid states that
+ don't have the expected keys or items, like a resource without connection or even a connection without path.
+ """
+ try:
+ yield
+ except KeyError as e:
+ logger.exception(
+ "The key `%s` could not be found in the DRBD configuration. The JSON format may have changed.", e
+ )
+ except Exception as e:
+ logger.exception("Failed to parse DRBD configuration: `%s`. The JSON format may have changed.", e)
+
+def _get_drbd_status(resource_name: str) -> Dict[str, Any]:
+ try:
+ stdout, stderr, ret_code = run_command([
+ _EXEC_PATH_DRBDSETUP, "status", resource_name, "--json"
+ ], simple=False)
+ if ret_code != 0:
+ logger.warning(
+ "Failed to get DRBD status of resource `%s`: `%s` (exit code %d).",
+ resource_name, stderr.strip(), ret_code
+ )
+ return {}
+ except Exception as e:
+ logger.error("Failed to get DRBD status of resource `%s`: `%s`.", resource_name, e)
+ return {}
+
+ try:
+ status = json.loads(stdout)
+ except Exception as e:
+ logger.error("Failed to read DRBD status of resource `%s` as JSON: `%s`.", resource_name, e)
+ return {}
+
+ with _handle_drbd_json_error():
+ return status[0]
+ return {}
+
+# ------------------------------------------------------------------------------
+
+@dataclass(frozen=True)
+class DrbdOpener:
+ pid: int
+ process_name: str
+ cmdline: List[str]
+ # The duration is expressed in milliseconds.
+ open_duration: int
+
+# ------------------------------------------------------------------------------
+
+class Drbd:
+ @staticmethod
+ def build_path(resource_name: str, volume_number: int) -> str:
+ _check_drbd_resource_name(resource_name)
+ _check_drbd_volume_number(volume_number)
+ return f"{DRBD_BY_RES_PATH}{resource_name}/{volume_number}"
+
+ @staticmethod
+ def get_name_from_path(path: str) -> str:
+ # Assume that we have a path like this:
+ # - "/dev/drbd/by-res//0"
+ # - "..//0"
+ if path.startswith(DRBD_BY_RES_PATH):
+ prefix_len = len(DRBD_BY_RES_PATH)
+ elif path.startswith("../"):
+ prefix_len = 3
+ else:
+ return ""
+
+ res_name_end = path.find("/", prefix_len)
+ if res_name_end == -1:
+ return ""
+
+ resource_name = path[prefix_len:res_name_end]
+ if _REGEX_DRBD_RESOURCE_NAME.fullmatch(resource_name) and path[res_name_end + 1:].isdecimal():
+ return resource_name
+ return ""
+
+ @staticmethod
+ def get_connection_address(resource_name: str, node_name: str) -> str:
+ _check_drbd_resource_name(resource_name)
+ status = _get_drbd_status(resource_name)
+ if not status:
+ return ""
+
+ with _handle_drbd_json_error():
+ for connection in status["connections"]:
+ if connection["name"] == node_name:
+ return connection["paths"][0]["remote_host"]["address"]
+ return ""
+
+ @staticmethod
+ def get_primary_address(resource_name: str) -> str:
+ _check_drbd_resource_name(resource_name)
+ status = _get_drbd_status(resource_name)
+ if not status:
+ return ""
+
+ with _handle_drbd_json_error():
+ if status["role"] == "Primary":
+ return status["connections"][0]["paths"][0]["this_host"]["address"]
+
+ for connection in status["connections"]:
+ if connection["peer-role"] == "Primary":
+ return connection["paths"][0]["remote_host"]["address"]
+
+ return ""
+
+ @staticmethod
+ def get_local_openers(resource_name: str, volume_number: int) -> List[DrbdOpener]:
+ _check_drbd_resource_name(resource_name)
+ _check_drbd_volume_number(volume_number)
+
+ path = Path(f"/sys/kernel/debug/drbd/resources/{resource_name}/volumes/{volume_number}/openers")
+ try:
+ lines = path.read_text().splitlines()
+ except Exception as e:
+ # The resource is probably available not on this node.
+ logger.info("Unable to get DRBD openers of volume `%s/%d`: `%s`.", resource_name, volume_number, e)
+ return []
+
+ drbd_openers = []
+ for line in lines:
+ match = _REGEX_DRBD_OPENER_LINE.fullmatch(line)
+ if not match:
+ logger.warning(
+ "Unable to parse DRBD opener line of volume `%s/%d` with: `%s`.",
+ resource_name,
+ volume_number,
+ line
+ )
+ continue
+
+ groups = match.groups()
+ pid = int(groups[1])
+ drbd_openers.append(DrbdOpener(
+ pid=pid,
+ process_name=groups[0],
+ # Note: `cmdline` is empty for `mount` calls. That's correct because `mount` process is dead.
+ cmdline=get_process_cmdline(pid),
+ open_duration=int(groups[2])
+ ))
+
+ return drbd_openers
+
+ @staticmethod
+ def demote(resource_name: str) -> bool:
+ _check_drbd_resource_name(resource_name)
+ error_message = ""
+ try:
+ _stdout, stderr, ret_code = run_command([_EXEC_PATH_DRBDSETUP, "secondary", resource_name], simple=False)
+ if not ret_code:
+ return True
+ error_message = stderr
+ except Exception as e:
+ error_message = str(e)
+ logger.error("Failed to demote DRBD resource `%s`: `%s`.", resource_name, error_message)
+ return False
diff --git a/src/xcp_storage/rpc/api.py b/src/xcp_storage/rpc/api.py
index 6344df5..7dabaa7 100644
--- a/src/xcp_storage/rpc/api.py
+++ b/src/xcp_storage/rpc/api.py
@@ -12,10 +12,12 @@
# You should have received a copy of the GNU General Public License
# along with this program. If not, see .
+import xcp_storage.rpc.modules.drbd as drbd
import xcp_storage.rpc.modules.echo as echo
# ==============================================================================
__all__ = [
+ "drbd",
"echo"
]
diff --git a/src/xcp_storage/rpc/modules/drbd.py b/src/xcp_storage/rpc/modules/drbd.py
new file mode 100644
index 0000000..ccadf2e
--- /dev/null
+++ b/src/xcp_storage/rpc/modules/drbd.py
@@ -0,0 +1,40 @@
+# Copyright (C) 2026 Vates SAS
+#
+# This program is free software: you can redistribute it and/or modify
+# it under the terms of the GNU General Public License as published by
+# the Free Software Foundation, either version 3 of the License, or
+# (at your option) any later version.
+# This program is distributed in the hope that it will be useful,
+# but WITHOUT ANY WARRANTY; without even the implied warranty of
+# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
+# GNU General Public License for more details.
+#
+# You should have received a copy of the GNU General Public License
+# along with this program. If not, see .
+
+from dataclasses import asdict
+
+from xcp_storage.backends.drbd import Drbd
+from xcp_storage.rpc.dispatcher import ApiDispatcher
+from xcp_storage.utils.json import JsonDict
+
+from xcp_storage.typing import (
+ assert_type,
+ cast,
+ Dict,
+ List,
+ MYPY,
+ PYREFLY,
+ Union,
+)
+
+# ==============================================================================
+
+@ApiDispatcher.method
+def get_openers(resource_name: str, volume_number: int) -> List[JsonDict]:
+ openers = [asdict(opener) for opener in Drbd.get_local_openers(resource_name, volume_number)]
+ # Some linters, such as pyrefly, may struggle to convert a type to a recursive type.
+ # We therefore perform a static check (so that it's validated by linters) and we cast it explicitly.
+ if PYREFLY and not MYPY:
+ assert_type(openers, List[Dict[str, Union[int, str, List[str]]]])
+ return cast(List[JsonDict], openers)
diff --git a/src/xcp_storage/typing/__init__.py b/src/xcp_storage/typing/__init__.py
index 4d944be..33284c9 100644
--- a/src/xcp_storage/typing/__init__.py
+++ b/src/xcp_storage/typing/__init__.py
@@ -63,3 +63,13 @@ def __call__(self, *_args, **_kwargs): # type: ignore # noqa: ANN002, ANN003, A
ParamSpec = _SubscriptableListMock() # type: ignore
Concatenate = _SubscriptableListMock() # type: ignore
+
+ if not hasattr(typing, "assert_type"):
+ _T = TypeVar("_T") # noqa: F405
+ def assert_type(value: _T, expected_type: Type[_T]) -> None: # noqa: F405
+ pass
+
+# ------------------------------------------------------------------------------
+
+MYPY: Final = False # noqa: F405
+PYREFLY: Final = False # noqa: F405
diff --git a/src/xcp_storage/utils/process.py b/src/xcp_storage/utils/process.py
index d7aa812..6a0be65 100644
--- a/src/xcp_storage/utils/process.py
+++ b/src/xcp_storage/utils/process.py
@@ -12,6 +12,7 @@
# You should have received a copy of the GNU General Public License
# along with this program. If not, see .
+from pathlib import Path
import subprocess
import xcp_storage.log as log
@@ -157,3 +158,13 @@ def run_command(
ret_code_callback=ret_code_callback,
quiet=quiet
)
+
+# ------------------------------------------------------------------------------
+
+def get_process_cmdline(pid: int) -> List[str]:
+ path = Path(f"/proc/{pid}/cmdline")
+ try:
+ return [arg.decode(errors="replace") for arg in path.read_bytes().split(b"\0") if arg]
+ except Exception as e:
+ logger.info("Unable to get command line of PID `%d`: `%s`.", pid, e)
+ return []
diff --git a/tests/backends/conftest.py b/tests/backends/conftest.py
new file mode 100644
index 0000000..dedf0e6
--- /dev/null
+++ b/tests/backends/conftest.py
@@ -0,0 +1,155 @@
+# Copyright (C) 2026 Vates SAS
+#
+# This program is free software: you can redistribute it and/or modify
+# it under the terms of the GNU General Public License as published by
+# the Free Software Foundation, either version 3 of the License, or
+# (at your option) any later version.
+# This program is distributed in the hope that it will be useful,
+# but WITHOUT ANY WARRANTY; without even the implied warranty of
+# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
+# GNU General Public License for more details.
+#
+# You should have received a copy of the GNU General Public License
+# along with this program. If not, see .
+
+import json
+
+import pytest
+
+# ==============================================================================
+
+@pytest.fixture
+def drbd_json_status_primary() -> str:
+ return json.dumps([{
+ "name": "xcp-volume-patate",
+ "node-id": 1,
+ "role": "Primary",
+ "suspended": False,
+ "suspended-user": False,
+ "suspended-no-data": False,
+ "suspended-fencing": False,
+ "suspended-quorum": False,
+ "force-io-failures": False,
+ "write-ordering": "flush",
+ "devices": [{
+ "volume": 0,
+ "minor": 1261,
+ "disk-state": "UpToDate",
+ "client": False,
+ "open": False,
+ "quorum": True,
+ "size": 2109208,
+ "read": 0,
+ "written": 0,
+ "al-writes": 0,
+ "bm-writes": 0,
+ "upper-pending": 0,
+ "lower-pending": 0
+ }],
+ "connections": [{
+ "peer-node-id": 0,
+ "name": "sr123-s1",
+ "connection-state": "Connected",
+ "congested": False,
+ "peer-role": "Secondary",
+ "tls": False,
+ "ap-in-flight": 0,
+ "rs-in-flight": 0,
+ "paths": [{
+ "this_host": {
+ "address": "10.10.0.13",
+ "port": 7261,
+ "family": "ipv4"
+ },
+ "remote_host": {
+ "address": "10.10.0.12",
+ "port": 7261,
+ "family": "ipv4"
+ },
+ "established": True
+ }],
+ "peer_devices": [{
+ "volume": 0,
+ "replication-state": "Established",
+ "peer-disk-state": "UpToDate",
+ "peer-client": False,
+ "resync-suspended": "no",
+ "received": 0,
+ "sent": 0,
+ "out-of-sync": 0,
+ "pending": 0,
+ "unacked": 0,
+ "has-sync-details": False,
+ "has-online-verify-details": False,
+ "percent-in-sync": 100
+ }]
+ }]
+ }])
+
+@pytest.fixture
+def drbd_json_status_secondary() -> str:
+ return json.dumps([{
+ "name": "xcp-volume-patate",
+ "node-id": 0,
+ "role": "Secondary",
+ "suspended": False,
+ "suspended-user": False,
+ "suspended-no-data": False,
+ "suspended-fencing": False,
+ "suspended-quorum": False,
+ "force-io-failures": False,
+ "write-ordering": "flush",
+ "devices": [{
+ "volume": 0,
+ "minor": 1261,
+ "disk-state": "UpToDate",
+ "client": False,
+ "open": False,
+ "quorum": True,
+ "size": 2109208,
+ "read": 0,
+ "written": 0,
+ "al-writes": 0,
+ "bm-writes": 0,
+ "upper-pending": 0,
+ "lower-pending": 0
+ }],
+ "connections": [{
+ "peer-node-id": 1,
+ "name": "sr123-s2",
+ "connection-state": "Connected",
+ "congested": False,
+ "peer-role": "Primary",
+ "tls": False,
+ "ap-in-flight": 0,
+ "rs-in-flight": 0,
+ "paths": [{
+ "this_host": {
+ "address": "10.10.0.12",
+ "port": 7261,
+ "family": "ipv4"
+ },
+ "remote_host": {
+ "address": "10.10.0.13",
+ "port": 7261,
+ "family": "ipv4"
+ },
+ "established": True
+ }],
+ "peer_devices": [{
+ "volume": 0,
+ "replication-state": "Established",
+ "peer-disk-state": "UpToDate",
+ "peer-client": False,
+ "resync-suspended": "no",
+ "received": 0,
+ "sent": 0,
+ "out-of-sync": 0,
+ "pending": 0,
+ "unacked": 0,
+ "has-sync-details": False,
+ "has-online-verify-details": False,
+ "percent-in-sync": 100
+ }]
+ }]
+ }])
diff --git a/tests/backends/test_drbd.py b/tests/backends/test_drbd.py
new file mode 100644
index 0000000..9bffe0f
--- /dev/null
+++ b/tests/backends/test_drbd.py
@@ -0,0 +1,336 @@
+# Copyright (C) 2026 Vates SAS
+#
+# This program is free software: you can redistribute it and/or modify
+# it under the terms of the GNU General Public License as published by
+# the Free Software Foundation, either version 3 of the License, or
+# (at your option) any later version.
+# This program is distributed in the hope that it will be useful,
+# but WITHOUT ANY WARRANTY; without even the implied warranty of
+# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
+# GNU General Public License for more details.
+#
+# You should have received a copy of the GNU General Public License
+# along with this program. If not, see .
+
+from pathlib import Path
+from unittest.mock import (
+ MagicMock,
+ patch,
+)
+
+import pytest
+
+from xcp_storage.backends.drbd import (
+ _get_drbd_status,
+ Drbd,
+ DrbdOpener,
+)
+
+from xcp_storage.typing import (
+ Any,
+ Dict,
+ Final,
+ List,
+)
+
+# ==============================================================================
+
+@pytest.mark.parametrize("resource_name, volume_number, expected_path", [
+ ("volume", 0, "/dev/drbd/by-res/volume/0"),
+ ("res-1", 2, "/dev/drbd/by-res/res-1/2"),
+ ("volume-2", 67, "/dev/drbd/by-res/volume-2/67")
+])
+def test_build_drbd_path(resource_name: str, volume_number: int, expected_path: str) -> None:
+ assert Drbd.build_path(resource_name, volume_number) == expected_path
+
+# ------------------------------------------------------------------------------
+
+@pytest.mark.parametrize("path, expected_name", [
+ ("/dev/drbd/by-res/xcp-volume-1/0", "xcp-volume-1"),
+ ("../xcp-volume-2/1", "xcp-volume-2"),
+ ("/dev/drbd/by-res/res-1/12", "res-1"),
+ ("/dev/drbd/by-res/xcp-volume-1/0/extra", ""),
+ ("/dev/drbd/by-res/a/b/c", ""),
+ ("../a/b/c", ""),
+ ("/dev/drbd/by-res/a/", ""),
+ ("../a/", ""),
+ ("/dev/drbd/by-res/a/x", ""),
+ ("/dev/drbd/by-res/a/-1", ""),
+ ("/dev/drbd/by-res/a/0/", ""),
+ ("../a/0/", ""),
+ ("/invalid/path/format", ""),
+ ("/dev/drbd/by-res/missing-volume-number", ""),
+ ("../missing-volume-number", ""),
+ ("", ""),
+ ("/dev/drbd/by-res//0", ""),
+ ("..//0", ""),
+ ("/dev/drbd/by-res/../0", ""),
+ ("../../x/0", ""),
+ ("../../../x/0", ""),
+ ("../../", ""),
+ ("/dev/drbd/by-res/-x/0", ""),
+ ("/dev/drbd/by-res/.ext/0", ""),
+ ("/dev/drbd/by-res/a b/0", "")
+])
+def test_get_drbd_name_from_path(path: str, expected_name: str) -> None:
+ assert Drbd.get_name_from_path(path) == expected_name
+
+# ------------------------------------------------------------------------------
+
+@patch("xcp_storage.backends.drbd.run_command")
+class TestGetDrbdStatus:
+ def test_success(self, mock_run_command: MagicMock, drbd_json_status_primary: str) -> None:
+ mock_run_command.return_value = (drbd_json_status_primary, "", 0)
+
+ status = _get_drbd_status("xcp-volume-patate")
+
+ assert status["name"] == "xcp-volume-patate"
+ assert status["role"] == "Primary"
+ mock_run_command.assert_called_once_with(
+ ["/usr/sbin/drbdsetup", "status", "xcp-volume-patate", "--json"], simple=False
+ )
+
+ def test_command_failure(self, mock_run_command: MagicMock, caplog: pytest.LogCaptureFixture) -> None:
+ mock_run_command.return_value = ("[]", "xcp-volume-frite: No such resource\n", 10)
+ with caplog.at_level("INFO"):
+ assert _get_drbd_status("xcp-volume-frite") == {}
+ assert (
+ "Failed to get DRBD status of resource `xcp-volume-frite`: "
+ "`xcp-volume-frite: No such resource` (exit code 10)."
+ ) in caplog.text
+
+ def test_command_not_found(self, mock_run_command: MagicMock, caplog: pytest.LogCaptureFixture) -> None:
+ mock_run_command.side_effect = FileNotFoundError(2, "No such file or directory", "/usr/sbin/drbdsetup")
+ with caplog.at_level("INFO"):
+ assert _get_drbd_status("xcp-volume-patate") == {}
+ assert (
+ "Failed to get DRBD status of resource `xcp-volume-patate`: "
+ "`[Errno 2] No such file or directory: '/usr/sbin/drbdsetup'`."
+ ) in caplog.text
+
+ @pytest.mark.parametrize("stdout", ["", "["])
+ def test_malformed_json(
+ self, mock_run_command: MagicMock, caplog: pytest.LogCaptureFixture, stdout: str
+ ) -> None:
+ mock_run_command.return_value = (stdout, "", 0)
+ with caplog.at_level("INFO"):
+ assert _get_drbd_status("xcp-volume-patate") == {}
+ assert "Failed to read DRBD status of resource `xcp-volume-patate` as JSON" in caplog.text
+
+ def test_empty_list(self, mock_run_command: MagicMock, caplog: pytest.LogCaptureFixture) -> None:
+ mock_run_command.return_value = ("[]", "", 0)
+ with caplog.at_level("INFO"):
+ assert _get_drbd_status("xcp-volume-patate") == {}
+ assert "Failed to parse DRBD configuration" in caplog.text
+
+# ------------------------------------------------------------------------------
+
+@patch("xcp_storage.backends.drbd.run_command")
+class TestGetDrbdConnectionAddress:
+ def test_connection_address_a(
+ self, mock_run_command: MagicMock, drbd_json_status_primary: str
+ ) -> None:
+ mock_run_command.return_value = (drbd_json_status_primary, "", 0)
+ assert Drbd.get_connection_address("xcp-volume-patate", "sr123-s1") == "10.10.0.12"
+
+ def test_connection_address_b(
+ self, mock_run_command: MagicMock, drbd_json_status_secondary: str
+ ) -> None:
+ mock_run_command.return_value = (drbd_json_status_secondary, "", 0)
+ assert Drbd.get_connection_address("xcp-volume-patate", "sr123-s2") == "10.10.0.13"
+
+ def test_unknown_node(self, mock_run_command: MagicMock, drbd_json_status_primary: str) -> None:
+ mock_run_command.return_value = (drbd_json_status_primary, "", 0)
+ assert Drbd.get_connection_address("xcp-volume-patate", "sr123-s67") == ""
+
+ def test_status_error(self, mock_run_command: MagicMock) -> None:
+ mock_run_command.return_value = ("", "xcp-volume-frite: No such resource\n", 10)
+ assert Drbd.get_connection_address("xcp-volume-frite", "sr123-s1") == ""
+
+# ------------------------------------------------------------------------------
+
+@patch("xcp_storage.backends.drbd.run_command")
+class TestGetDrbdPrimaryAddress:
+ def test_primary_local(
+ self, mock_run_command: MagicMock, drbd_json_status_primary: str
+ ) -> None:
+ mock_run_command.return_value = (drbd_json_status_primary, "", 0)
+ assert Drbd.get_primary_address("xcp-volume-patate") == "10.10.0.13"
+
+ def test_primary_remote(
+ self, mock_run_command: MagicMock, drbd_json_status_secondary: str
+ ) -> None:
+ mock_run_command.return_value = (drbd_json_status_secondary, "", 0)
+ assert Drbd.get_primary_address("xcp-volume-patate") == "10.10.0.13"
+
+ def test_status_error(self, mock_run_command: MagicMock) -> None:
+ mock_run_command.return_value = ("", "xcp-volume-frite: No such resource\n", 10)
+ assert Drbd.get_primary_address("xcp-volume-frite") == ""
+
+# ------------------------------------------------------------------------------
+
+@patch.object(Path, "read_text")
+class TestGetDrbdLocalOpeners:
+ @patch("xcp_storage.backends.drbd.get_process_cmdline")
+ def test_multi_openers(self, mock_get_cmdline: MagicMock, mock_read_text: MagicMock) -> None:
+ data: List[Dict[str, Any]] = [{
+ "pid": 482584,
+ "process_name": "tapback",
+ "cmdline": ["tapback", "-d", "-x", "1"],
+ "open_duration": 86143
+ }, {
+ "invalid_cmdline": "invalid_line\n"
+ }, {
+ "pid": 483388,
+ "process_name": "python",
+ "cmdline": ["storage"],
+ "open_duration": 1877
+ }]
+ opener_data = [d for d in data if not d.get("invalid_cmdline")]
+
+ mock_read_text.return_value = "".join(
+ f"{d['process_name']} {d['pid']} {d['open_duration']}\n"
+ if not d.get("invalid_cmdline")
+ else d["invalid_cmdline"]
+ for d in data
+ )
+
+ cmdline_mapping = {d["pid"]: d["cmdline"] for d in opener_data}
+ mock_get_cmdline.side_effect = lambda pid: cmdline_mapping.get(pid, [])
+
+ openers = Drbd.get_local_openers("res-test", 0)
+
+ assert len(openers) == len(opener_data)
+ for i, expected in enumerate(opener_data):
+ assert openers[i] == DrbdOpener(
+ pid=expected["pid"],
+ process_name=expected["process_name"],
+ cmdline=expected["cmdline"],
+ open_duration=expected["open_duration"],
+ )
+
+ @pytest.mark.parametrize("line", [
+ "python 483388 1877 extra",
+ "python 483388",
+ "python",
+ "483388 1877",
+ ""
+ ])
+ def test_invalid_opener_line_is_ignored(
+ self, mock_read_text: MagicMock, caplog: pytest.LogCaptureFixture, line: str
+ ) -> None:
+ mock_read_text.return_value = f"{line}\n"
+
+ with caplog.at_level("WARNING"):
+ assert Drbd.get_local_openers("res-test", 0) == []
+ assert f"Unable to parse DRBD opener line of volume `res-test/0` with: `{line}`." in caplog.text
+
+ def test_file_not_found(
+ self, mock_read_text: MagicMock, caplog: pytest.LogCaptureFixture
+ ) -> None:
+ openers_path = "/sys/kernel/debug/drbd/resources/res-test/volumes/0/openers"
+ mock_read_text.side_effect = FileNotFoundError(
+ 2,
+ "No such file or directory",
+ openers_path
+ )
+ with caplog.at_level("INFO"):
+ assert Drbd.get_local_openers("res-test", 0) == []
+ assert (
+ "Unable to get DRBD openers of volume `res-test/0`: "
+ f"`[Errno 2] No such file or directory: '{openers_path}'`."
+ ) in caplog.text
+
+# ------------------------------------------------------------------------------
+
+@patch("xcp_storage.backends.drbd.run_command")
+class TestDemoteDrbd:
+ def test_success(self, mock_run_command: MagicMock) -> None:
+ mock_run_command.return_value = ("", "", 0)
+ assert Drbd.demote("res-test")
+
+ def test_command_not_found(self, mock_run_command: MagicMock, caplog: pytest.LogCaptureFixture) -> None:
+ mock_run_command.side_effect = FileNotFoundError(2, "No such file or directory", "/usr/sbin/drbdsetup")
+ with caplog.at_level("INFO"):
+ assert not Drbd.demote("res-test")
+ assert (
+ "Failed to demote DRBD resource `res-test`: "
+ "`[Errno 2] No such file or directory: '/usr/sbin/drbdsetup'`."
+ ) in caplog.text
+
+ def test_demote_drbd_open_on_another_node(
+ self, mock_run_command: MagicMock, caplog: pytest.LogCaptureFixture
+ ) -> None:
+ stderr = (
+ "xcp-volume-ad084d58-c12b-42df-affa-8cda2321580b: State change failed: "
+ "(-12) Device is held open by someone\n"
+ "additional info from kernel:\n"
+ "/dev/drbd1261 open_ro_cnt:0, open_rw_cnt:2; list of openers follows\n"
+ "drbd1261 opened by python (pid 483388) at 2026-06-30 21:15:38.690\n"
+ "drbd1261 opened by python (pid 482584) at 2026-06-30 21:14:14.425\n"
+ )
+ mock_run_command.return_value = ("", stderr, 11)
+
+ with caplog.at_level("INFO"):
+ assert not Drbd.demote("res-test")
+ assert f"Failed to demote DRBD resource `res-test`: `{stderr}`." in caplog.text
+
+# ------------------------------------------------------------------------------
+
+# Resource names and volume numbers can come from an RPC call or external stack: types are not checked at runtime.
+# So we test different situations.
+@patch.object(Path, "read_text")
+@patch("xcp_storage.backends.drbd.run_command")
+class TestInvalidDrbdArguments:
+ INVALID_RESOURCE_NAMES: Final = (
+ "", ".", "..", "../x", "x/../y", "a/b", "/abs", "a b", "a\n", "a\x00b", "\u00e9", "-x", "--param", ".ext"
+ )
+ NOT_RESOURCE_NAMES: Final = (None, 1, b"res")
+
+ INVALID_VOLUME_NUMBERS: Final = (-1, -42)
+ NOT_VOLUME_NUMBERS: Final = ("0", "0/../../x", 1.5, True, False, None)
+
+ @pytest.mark.parametrize("resource_name, message", [
+ *((name, "Invalid DRBD resource name") for name in INVALID_RESOURCE_NAMES),
+ *((name, "Not a DRBD resource name") for name in NOT_RESOURCE_NAMES)
+ ])
+ def test_invalid_resource_name(
+ self,
+ mock_run_command: MagicMock,
+ mock_read_text: MagicMock,
+ resource_name: Any, # noqa: ANN401
+ message: str
+ ) -> None:
+ for call in (
+ lambda: Drbd.build_path(resource_name, 0),
+ lambda: Drbd.get_connection_address(resource_name, "sr123-s1"),
+ lambda: Drbd.get_primary_address(resource_name),
+ lambda: Drbd.get_local_openers(resource_name, 0),
+ lambda: Drbd.demote(resource_name)
+ ):
+ with pytest.raises(ValueError, match=message):
+ call()
+
+ mock_run_command.assert_not_called()
+ mock_read_text.assert_not_called()
+
+ @pytest.mark.parametrize("volume_number, message", [
+ *((number, "Invalid DRBD volume number") for number in INVALID_VOLUME_NUMBERS),
+ *((number, "Not a DRBD volume number") for number in NOT_VOLUME_NUMBERS)
+ ])
+ def test_invalid_volume_number(
+ self,
+ mock_run_command: MagicMock,
+ mock_read_text: MagicMock,
+ volume_number: Any, # noqa: ANN401
+ message: str
+ ) -> None:
+ for call in (
+ lambda: Drbd.build_path("res-test", volume_number),
+ lambda: Drbd.get_local_openers("res-test", volume_number)
+ ):
+ with pytest.raises(ValueError, match=message):
+ call()
+
+ mock_run_command.assert_not_called()
+ mock_read_text.assert_not_called()
diff --git a/tests/utils/test_process.py b/tests/utils/test_process.py
new file mode 100644
index 0000000..e896ee6
--- /dev/null
+++ b/tests/utils/test_process.py
@@ -0,0 +1,52 @@
+# Copyright (C) 2026 Vates SAS
+#
+# This program is free software: you can redistribute it and/or modify
+# it under the terms of the GNU General Public License as published by
+# the Free Software Foundation, either version 3 of the License, or
+# (at your option) any later version.
+# This program is distributed in the hope that it will be useful,
+# but WITHOUT ANY WARRANTY; without even the implied warranty of
+# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
+# GNU General Public License for more details.
+#
+# You should have received a copy of the GNU General Public License
+# along with this program. If not, see .
+
+from pathlib import Path
+from unittest.mock import MagicMock, patch
+
+import pytest
+
+from xcp_storage.utils.process import get_process_cmdline
+
+from xcp_storage.typing import Final
+
+# ==============================================================================
+
+@patch.object(Path, "read_bytes")
+class TestGetProcessCmdline:
+ PID: Final = 67
+
+ def test_success(self, mock_read_bytes: MagicMock) -> None:
+ mock_read_bytes.return_value = b"python\x00-m\x00pytest\x00"
+ assert get_process_cmdline(self.PID) == ["python", "-m", "pytest"]
+
+ def test_null_args(self, mock_read_bytes: MagicMock) -> None:
+ mock_read_bytes.return_value = b"abricot\x00\x00\x00-orange\x00\x00jus ananas \x00"
+ assert get_process_cmdline(self.PID) == ["abricot", "-orange", "jus ananas "]
+
+ def test_non_utf8_args(self, mock_read_bytes: MagicMock) -> None:
+ mock_read_bytes.return_value = b"prog\x00\xff\xfe\x00arg\x00"
+ assert get_process_cmdline(self.PID) == ["prog", "\ufffd\ufffd", "arg"]
+
+ def test_pid_not_found(self, mock_read_bytes: MagicMock, caplog: pytest.LogCaptureFixture) -> None:
+ cmdline_path = f"/proc/{self.PID}/cmdline"
+ mock_read_bytes.side_effect = FileNotFoundError(2, "No such file or directory", cmdline_path)
+
+ with caplog.at_level("INFO"):
+ assert get_process_cmdline(self.PID) == []
+
+ assert (
+ f"Unable to get command line of PID `{self.PID}`: "
+ f"`[Errno 2] No such file or directory: '{cmdline_path}'`."
+ ) in caplog.text