Skip to content
Open
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
Empty file.
72 changes: 72 additions & 0 deletions src/xcp_storage/backends/linstor/controller.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
# 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 <https://www.gnu.org/licenses/>.

from xcp_storage.backends.linstor.satellite import LINSTOR_SATELLITE_PORT_PLAIN, LINSTOR_SATELLITE_PORT_SSL
from xcp_storage.config.platform import get_exec_path
from xcp_storage.utils.process import run_command
from xcp_storage.utils.service import (
is_service_active,
restart_service,
start_service,
stop_service,
)

from xcp_storage.typing import Final, List

# ==============================================================================

LINSTOR_CONTROLLER_PORT_PLAIN: Final = 3370
LINSTOR_CONTROLLER_PORT_SSL: Final = 3371

# ------------------------------------------------------------------------------

_EXEC_PATH_SS: Final = get_exec_path("/usr/sbin/ss", {"debian": "/usr/bin/ss"})

_SERVICE_LINSTOR_CONTROLLER: Final = "linstor-controller"

# ------------------------------------------------------------------------------

class LinstorController:
@staticmethod
def get_addresses() -> List[str]:
stdout = run_command([
_EXEC_PATH_SS, "-tnpH", "state", "established",
f"( sport = :{LINSTOR_SATELLITE_PORT_PLAIN} or sport = :{LINSTOR_SATELLITE_PORT_SSL} )"
], expected_ret_code=0)
return [
line.split()[3].rsplit(":", 1)[0]
for line in stdout.splitlines()
]

@classmethod
def get_uri(cls) -> str:
# TODO(XCPNG-3033): On caller side, check that an IP address from the current pool is returned.
addresses = cls.get_addresses()
return "linstor://" + addresses[0] if addresses else ""

@staticmethod
def is_running() -> bool:
return is_service_active(_SERVICE_LINSTOR_CONTROLLER)

@staticmethod
def start() -> None:
start_service(_SERVICE_LINSTOR_CONTROLLER)

@staticmethod
def stop() -> None:
stop_service(_SERVICE_LINSTOR_CONTROLLER)

@staticmethod
def restart() -> None:
restart_service(_SERVICE_LINSTOR_CONTROLLER)
45 changes: 45 additions & 0 deletions src/xcp_storage/backends/linstor/satellite.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
# 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 <https://www.gnu.org/licenses/>.

from xcp_storage.utils.service import (
disable_and_stop_service,
enable_and_start_service,
is_service_active,
)

from xcp_storage.typing import Final

# ==============================================================================

LINSTOR_SATELLITE_PORT_PLAIN: Final = 3366
LINSTOR_SATELLITE_PORT_SSL: Final = 3367

# ------------------------------------------------------------------------------

_SERVICE_LINSTOR_SATELLITE: Final = "linstor-satellite"

# ------------------------------------------------------------------------------

class LinstorSatellite:
@staticmethod
def is_running() -> bool:
return is_service_active(_SERVICE_LINSTOR_SATELLITE)

@staticmethod
def enable_and_start() -> None:
enable_and_start_service(_SERVICE_LINSTOR_SATELLITE)

@staticmethod
def disable_and_stop() -> None:
disable_and_stop_service(_SERVICE_LINSTOR_SATELLITE)
37 changes: 36 additions & 1 deletion src/xcp_storage/config/platform.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,10 @@
# You should have received a copy of the GNU General Public License
# along with this program. If not, see <https://www.gnu.org/licenses/>.

from xcp_storage.typing import Final
from functools import lru_cache
from pathlib import Path

from xcp_storage.typing import Final, Mapping, Tuple

# ==============================================================================
# Attributes that depend on the execution environment.
Expand All @@ -21,3 +24,35 @@
# ==============================================================================

DEFAULT_FIREWALL_INPUT_CHAIN: Final = "xapi-INPUT"

# ------------------------------------------------------------------------------

_OS_RELEASE_PATH: Final = "/etc/os-release"

@lru_cache(maxsize=None)
def get_os_ids() -> Tuple[str, ...]:
"""
Get the IDs of the current distribution, most specific first: `ID` followed
by the values of `ID_LIKE` (e.g. `("ubuntu", "debian")`). The result is
cached. An empty tuple is returned if `/etc/os-release` can't be read.
"""

try:
lines = Path(_OS_RELEASE_PATH).read_text(encoding="utf-8").splitlines()
except OSError:
return ()

values = {}
for line in lines:
key, _separator, value = line.partition("=")
values[key.strip()] = value.strip().strip("\"'")

return tuple(values.get("ID", "").split() + values.get("ID_LIKE", "").split())

def get_exec_path(default: str, by_os_id: Mapping[str, str]) -> str:
"""
Get the path of an executable: the one registered for the first matching OS
ID, otherwise `default`.
"""

return next((by_os_id[os_id] for os_id in get_os_ids() if os_id in by_os_id), default)
17 changes: 16 additions & 1 deletion src/xcp_storage/network/socket.py
Original file line number Diff line number Diff line change
Expand Up @@ -198,13 +198,18 @@ def create_client_sock(
reuse_address: bool = True,
keep_alive: bool = True,
timeout: Optional[float] = None,
ssl_context: Optional[ssl.SSLContext] = None
ssl_context: Optional[ssl.SSLContext] = None,
source_address: Optional[str] = None,
source_port: int = 0
) -> socket.socket:
family, connect = format_address(address, port)
sock = _create_stream_sock(address, family, bind=False, reuse_address=reuse_address, ssl_context=ssl_context)
_normalize_and_set_sock_timeout(sock, timeout)

try:
if source_address:
_, bind = format_address(source_address, source_port)
sock.bind(bind)
sock.connect(connect)
except OSError as e:
with contextlib.suppress(Exception):
Expand Down Expand Up @@ -318,6 +323,12 @@ def socket_wait_readable(sock: socket.socket, *, timeout: Optional[float] = None
def get_socket_family_str(sock: socket.socket) -> str:
return _FAMILY_TO_STR.get(sock.family, "Unknown")

def get_socket_address(sock: socket.socket) -> Optional[str]:
try:
return sock.getsockname()[0]
except OSError:
return None

def get_socket_port(sock: socket.socket) -> Optional[int]:
try:
return sock.getsockname()[1]
Expand Down Expand Up @@ -371,6 +382,10 @@ def close(self) -> None:
def family_str(self) -> str:
return get_socket_family_str(self.sock)

@property
def address(self) -> Optional[str]:
return get_socket_address(self.sock)

@property
def port(self) -> Optional[int]:
return get_socket_port(self.sock)
Expand Down
Loading
Loading