mirror of
https://github.com/MacRimi/ProxMenux.git
synced 2026-10-08 22:46:41 +00:00
feat(oci): recover applications after a Proxmox reinstall, cluster records and AMD GPU profiles
The installation record of an OCI application travels with its container: a copy inside the container and another in /etc/pve, written together with the one kept on the host. A container restored on a newly installed Proxmox, restored with another ID or moved to another node of a cluster is recognised and registered again, with its private network, hookscript, Rclone mount, host firewall rule and NVIDIA runtime. The Monitor offers the same recovery from the Updates tab. AMD GPUs are offered by generation. A GPU the ROCm image supports takes the profile as it is; one of a supported family (Radeon 680M, 780M) is an experimental option that asks for confirmation and is never proposed; an older one is not offered. The GPU is checked with a real inference before the installation accepts it. Recreate changes what runs recognition in an installed Immich, between the CPU and a GPU of the host. Updates: - A failed update that is restored and checked removes its temporary container and the disks of the failed attempt. - Every container volume is part of the backups, so Jellyfin, Plex and Hugo update with their default installation. - An image published with a Docker-format manifest is recognised by its layers and build time and updates. - The Proxmox notes of a multi-container application link to its LAN address.
This commit is contained in:
@@ -176,6 +176,12 @@ oci_quiet pct unmount "$VMID"
|
||||
sed -i -E '/^lxc\.environment\.runtime: RCLONE_RC_(USER|PASS)=/d' "/etc/pve/lxc/${VMID}.conf"
|
||||
oci_quiet pct set "$VMID" --entrypoint /usr/local/bin/rclone-mount-lxc-start
|
||||
sed -i -E '/^hookscript:/d' "/etc/pve/lxc/${VMID}.conf"
|
||||
# A new Proxmox installation accepts no snippets on `local`.
|
||||
if ! pvesm status --content snippets 2>/dev/null | awk 'NR > 1 && $1 == "local" {found=1} END {exit !found}'; then
|
||||
LOCAL_CONTENT=$(pvesh get /storage/local --output-format json 2>/dev/null | jq -r '.content // empty')
|
||||
[[ -n $LOCAL_CONTENT ]] || die "$(translate "The local storage does not accept snippets")"
|
||||
oci_quiet pvesm set local --content "${LOCAL_CONTENT},snippets"
|
||||
fi
|
||||
oci_quiet pct set "$VMID" --hookscript "local:snippets/$(basename "$HOOK_PATH")"
|
||||
msg_ok "$(translate "Mount mode applied")"
|
||||
|
||||
|
||||
@@ -98,6 +98,9 @@ done
|
||||
[[ -r $STACK_DEPENDENCY_HOOK ]] || die "$(translate "The stack startup hook was not found")"
|
||||
|
||||
ML_ACCELERATION=$(jq -er '.machine_learning.acceleration // "cpu"' "$DEPLOYMENT_FILE")
|
||||
ML_GFX_OVERRIDE=$(jq -r '.machine_learning.gfx_override // empty' "$DEPLOYMENT_FILE")
|
||||
[[ -z $ML_GFX_OVERRIDE || ( $ML_ACCELERATION == rocm && $ML_GFX_OVERRIDE =~ ^[0-9]{1,2}\.[0-9]\.[0-9]$ ) ]] \
|
||||
|| die "$(translate "Invalid ROCm generation override")"
|
||||
source "$SCRIPT_DIR/oci_nvidia_setup.sh"
|
||||
source "$SCRIPT_DIR/oci_immich_ml.sh"
|
||||
validate_immich_ml_profile
|
||||
@@ -437,6 +440,10 @@ set_runtime_env "$ML_ID" MACHINE_LEARNING_CACHE_FOLDER /cache
|
||||
set_runtime_env "$ML_ID" TRANSFORMERS_CACHE /cache
|
||||
set_runtime_env "$ML_ID" MACHINE_LEARNING_MODEL_INTRA_OP_THREADS 2
|
||||
set_runtime_env "$ML_ID" MACHINE_LEARNING_MODEL_INTER_OP_THREADS 1
|
||||
if [[ -n $ML_GFX_OVERRIDE ]]; then
|
||||
set_runtime_env "$ML_ID" HSA_OVERRIDE_GFX_VERSION "$ML_GFX_OVERRIDE"
|
||||
set_runtime_env "$ML_ID" HSA_USE_SVM 0
|
||||
fi
|
||||
configure_immich_ml_gpu
|
||||
oci_quiet pct mount "$ML_ID"
|
||||
ML_ROOT="/var/lib/lxc/${ML_ID}/rootfs"
|
||||
|
||||
@@ -285,7 +285,8 @@ apply_host_monitor_firewall() {
|
||||
apply_host_monitor() {
|
||||
[[ -n ${HOST_MONITOR:-} ]] || return 0
|
||||
# PVE permits lxc.include but not namespace keys directly in the CT config.
|
||||
# This static, cluster-persistent companion must accompany cross-host restores.
|
||||
# This static, cluster-persistent companion is written again by the recovery
|
||||
# of a container restored on another host.
|
||||
# /etc/pve/proxmenux is the same on every node of a cluster; /etc/pve/lxc
|
||||
# is the folder of this node only.
|
||||
local include=/etc/pve/proxmenux/host-monitor native
|
||||
@@ -303,7 +304,7 @@ apply_host_monitor() {
|
||||
# Do not remove the Proxmox pre-start, autodev or post-stop hooks.
|
||||
set_lxc_directive lxc.hook.mount ""
|
||||
msg_ok "$(translate "Host monitor configured: shared PID and network namespaces, LXCFS disabled in this container")"
|
||||
msg_info2 "$(translate "Every node of this cluster already has this file. If you restore this container on any other Proxmox host, copy it to the same path first, because the container backup does not include it:") $include"
|
||||
msg_info2 "$(translate "Every node of this cluster already has this file. A backup of the container does not include it: on another Proxmox host it is written again when the application is recovered from Manage installed OCI applications:") $include"
|
||||
}
|
||||
|
||||
verify_host_monitor() {
|
||||
|
||||
@@ -0,0 +1,385 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Copies of the record of an OCI application that outlive the host record.
|
||||
|
||||
The record of an installation lives on the disk of the node that installed
|
||||
it. Two copies are kept with it, written together after every operation:
|
||||
|
||||
- inside the container, in its root filesystem, so a backup carries it: a
|
||||
container restored on another Proxmox host, or on this one after a
|
||||
reinstall, still has what is needed to register it again. The container
|
||||
could change this copy, so it is checked against the configuration Proxmox
|
||||
restored before anything is taken from it;
|
||||
- in /etc/pve, which every node of a cluster shares and only root of the host
|
||||
reads, so a container that migrates finds its record on the node it moves
|
||||
to. This one is trusted.
|
||||
|
||||
Neither is read back as a source of truth while the host record is current.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import contextlib
|
||||
import datetime
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
import re
|
||||
import socket
|
||||
import stat
|
||||
import subprocess
|
||||
import sys
|
||||
import uuid
|
||||
|
||||
import oci_instances as instances
|
||||
from oci_installation_state import sha
|
||||
from oci_ui import translate
|
||||
|
||||
DIRECTORY = '.proxmenux'
|
||||
NAME = 'oci-record.json'
|
||||
KIND = 'proxmenux.oci-carried-record'
|
||||
STAMP = 'carried.sha256'
|
||||
LIMIT = 16 * 1024 * 1024
|
||||
STACK_CONTRACT = '/etc/pve/priv/proxmenux-stack-{}.json'
|
||||
NVIDIA_HOOK = re.compile(r'/usr/local/lib/proxmenux/oci/nvidia-mount-([a-f0-9]{64})\.sh')
|
||||
SNIPPETS = Path('/var/lib/vz/snippets')
|
||||
# Every node of a cluster reads this folder and only root of the host can:
|
||||
# a container that moves to another node finds its record there.
|
||||
CLUSTER = Path('/etc/pve/priv/proxmenux/oci')
|
||||
RCLONE_HOOK = re.compile(r'^hookscript: local:snippets/(proxmenux-rclone-[0-9]+-fuse-hook\.sh)$', re.MULTILINE)
|
||||
|
||||
|
||||
def _run(*args):
|
||||
return subprocess.run(args, capture_output=True, text=True, timeout=120, check=False)
|
||||
|
||||
|
||||
def running_pid(vmid):
|
||||
result = _run('lxc-info', '-n', str(int(vmid)), '-pH')
|
||||
pid = result.stdout.strip()
|
||||
return int(pid) if result.returncode == 0 and pid.isdigit() else None
|
||||
|
||||
|
||||
@contextlib.contextmanager
|
||||
def container_root(vmid, mount_stopped=True):
|
||||
"""The root filesystem of the container as the host sees it, or None. A
|
||||
running container is reached through its init process; a stopped one is
|
||||
mounted for as long as the block lasts."""
|
||||
pid = running_pid(vmid)
|
||||
if pid:
|
||||
yield Path(f'/proc/{pid}/root')
|
||||
return
|
||||
if not mount_stopped or _run('pct', 'mount', str(int(vmid))).returncode != 0:
|
||||
yield None
|
||||
return
|
||||
try:
|
||||
yield Path(f'/var/lib/lxc/{int(vmid)}/rootfs')
|
||||
finally:
|
||||
_run('pct', 'unmount', str(int(vmid)))
|
||||
|
||||
|
||||
def mapped_root(config):
|
||||
"""The host owner that is root inside the container."""
|
||||
owner = {}
|
||||
for line in config.splitlines():
|
||||
fields = line.partition(': ')[2].split()
|
||||
if line.startswith('lxc.idmap: ') and len(fields) == 4 and fields[1] == '0' and fields[0] in 'ug':
|
||||
owner.setdefault(fields[0], int(fields[2]))
|
||||
default = 100000 if re.search(r'^unprivileged: 1$', config, re.MULTILINE) else 0
|
||||
return owner.get('u', default), owner.get('g', default)
|
||||
|
||||
|
||||
def _directory(root, owner=None):
|
||||
"""The private folder of the copy, opened without following a link the
|
||||
container could have left in its place."""
|
||||
root_fd = os.open(root, os.O_RDONLY | os.O_DIRECTORY)
|
||||
try:
|
||||
if owner is not None:
|
||||
with contextlib.suppress(FileExistsError):
|
||||
os.mkdir(DIRECTORY, 0o700, dir_fd=root_fd)
|
||||
fd = os.open(DIRECTORY, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW, dir_fd=root_fd)
|
||||
finally:
|
||||
os.close(root_fd)
|
||||
if owner is not None:
|
||||
os.fchown(fd, *owner)
|
||||
os.fchmod(fd, 0o700)
|
||||
return fd
|
||||
|
||||
|
||||
def valid_copy(value):
|
||||
return (isinstance(value, dict) and value.get('kind') == KIND and value.get('schema_version') == 1
|
||||
and isinstance(value.get('record'), dict))
|
||||
|
||||
|
||||
def read_cluster(vmid):
|
||||
"""The copy the cluster keeps for a container, or None."""
|
||||
path = CLUSTER / f'{int(vmid)}.json'
|
||||
try:
|
||||
if path.is_symlink() or path.stat().st_size > LIMIT:
|
||||
return None
|
||||
value = json.loads(path.read_text())
|
||||
except (OSError, ValueError):
|
||||
return None
|
||||
return value if valid_copy(value) else None
|
||||
|
||||
|
||||
def write_cluster(vmid, value):
|
||||
"""Best effort: /etc/pve is read-only without quorum and takes files of
|
||||
up to 1 MiB. The copy inside the container does not depend on it."""
|
||||
try:
|
||||
CLUSTER.mkdir(parents=True, exist_ok=True)
|
||||
(CLUSTER / f'{int(vmid)}.json').write_text(json.dumps(value))
|
||||
except OSError:
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
def remove_cluster(vmid):
|
||||
with contextlib.suppress(OSError):
|
||||
(CLUSTER / f'{int(vmid)}.json').unlink()
|
||||
|
||||
|
||||
def keep_cluster(vmid, value, generation):
|
||||
"""Leave in the cluster the same copy the container carries."""
|
||||
shared = read_cluster(vmid)
|
||||
if shared is not None and digest(shared) == digest(value) and shared.get('generation') == generation:
|
||||
return
|
||||
write_cluster(vmid, dict(value, generation=generation, saved_at=value.get('saved_at')
|
||||
or datetime.datetime.now(datetime.timezone.utc).isoformat()))
|
||||
|
||||
|
||||
def replaced_elsewhere(record, record_path, shared, node):
|
||||
"""Whether another node of the cluster wrote a different record for this
|
||||
container after the one of this host: the container was changed there and
|
||||
came back."""
|
||||
if shared is None or shared.get('node') == node or shared['record'] == record \
|
||||
or shared['record'].get('installation_id') != record['installation_id'] \
|
||||
or shared['record'].get('vmid') != record['vmid']:
|
||||
return False
|
||||
try:
|
||||
saved = datetime.datetime.fromisoformat(str(shared.get('saved_at')))
|
||||
return saved.timestamp() > record_path.stat().st_mtime
|
||||
except (ValueError, OSError, TypeError):
|
||||
return False
|
||||
|
||||
|
||||
def read_copy(root):
|
||||
"""The copy a container carries, or None when it has none that can be used."""
|
||||
try:
|
||||
directory = _directory(root)
|
||||
except OSError:
|
||||
return None
|
||||
try:
|
||||
fd = os.open(NAME, os.O_RDONLY | os.O_NOFOLLOW, dir_fd=directory)
|
||||
except OSError:
|
||||
return None
|
||||
finally:
|
||||
os.close(directory)
|
||||
with os.fdopen(fd, 'rb') as source:
|
||||
info = os.fstat(source.fileno())
|
||||
if not stat.S_ISREG(info.st_mode) or info.st_size > LIMIT:
|
||||
return None
|
||||
try:
|
||||
value = json.loads(source.read(LIMIT))
|
||||
except ValueError:
|
||||
return None
|
||||
return value if valid_copy(value) else None
|
||||
|
||||
|
||||
def write_copy(root, value, owner):
|
||||
directory = _directory(root, owner)
|
||||
temporary = f'.{NAME}.{os.getpid()}'
|
||||
try:
|
||||
with contextlib.suppress(FileNotFoundError):
|
||||
os.unlink(temporary, dir_fd=directory)
|
||||
fd = os.open(temporary, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, 0o600, dir_fd=directory)
|
||||
try:
|
||||
with os.fdopen(fd, 'w') as output:
|
||||
json.dump(value, output)
|
||||
output.flush()
|
||||
os.fchown(output.fileno(), *owner)
|
||||
os.fsync(output.fileno())
|
||||
os.rename(temporary, NAME, src_dir_fd=directory, dst_dir_fd=directory)
|
||||
except BaseException:
|
||||
with contextlib.suppress(OSError):
|
||||
os.unlink(temporary, dir_fd=directory)
|
||||
raise
|
||||
os.fsync(directory)
|
||||
finally:
|
||||
os.close(directory)
|
||||
|
||||
|
||||
def rclone_mount(config):
|
||||
"""Where the hookscript of an Rclone container publishes its mount on the
|
||||
host. The hookscript itself is written again from the engine."""
|
||||
hook = RCLONE_HOOK.search(config)
|
||||
if not hook:
|
||||
return None
|
||||
path = SNIPPETS / hook[1]
|
||||
try:
|
||||
if path.is_symlink():
|
||||
return None
|
||||
text = path.read_text()
|
||||
except (OSError, UnicodeDecodeError):
|
||||
return None
|
||||
name = re.search(r'^inside=/data/mounts/(\S+)$', text, re.MULTILINE)
|
||||
views = {key: re.search(rf'^{key}=(/\S+)/(\S+)$', text, re.MULTILINE) for key in ('published', 'published_ro')}
|
||||
parent = re.search(r'^\s*if ! mountpoint -q (/\S+); then$', text, re.MULTILINE)
|
||||
if not name or not parent or not all(views.values()) \
|
||||
or any(view[2] != name[1] for view in views.values()):
|
||||
return None
|
||||
return {'mount_name': name[1], 'shared_mount_root': views['published'][1],
|
||||
'shared_mount_read_only_root': views['published_ro'][1], 'shared_mount_root_parent': parent[1]}
|
||||
|
||||
|
||||
def bundle(record, config):
|
||||
"""What the container carries: its record and the host files a start of
|
||||
the application depends on."""
|
||||
vmid = record['vmid']
|
||||
value = {'schema_version': 1, 'kind': KIND, 'node': socket.gethostname().split('.', 1)[0],
|
||||
'vmid': vmid, 'record': record}
|
||||
contract = Path(STACK_CONTRACT.format(vmid))
|
||||
if contract.is_file() and not contract.is_symlink():
|
||||
with contextlib.suppress(OSError, ValueError):
|
||||
value['stack_contract'] = json.loads(contract.read_text())
|
||||
hook = NVIDIA_HOOK.search(config)
|
||||
if hook:
|
||||
path = Path(hook[0])
|
||||
with contextlib.suppress(OSError, UnicodeDecodeError):
|
||||
if path.is_file() and not path.is_symlink() and sha(path.read_bytes()) == hook[1]:
|
||||
value['nvidia_hook'] = {'sha256': hook[1], 'content': path.read_text()}
|
||||
mount = rclone_mount(config)
|
||||
if mount:
|
||||
value['rclone_mount'] = mount
|
||||
return value
|
||||
|
||||
|
||||
def digest(value):
|
||||
stable = {key: item for key, item in value.items() if key not in ('saved_at', 'generation')}
|
||||
return sha(json.dumps(stable, sort_keys=True, separators=(',', ':')).encode())
|
||||
|
||||
|
||||
def read_stamp(path):
|
||||
"""What this host last wrote into the container: the digest of the copy
|
||||
and the mark that tells that copy from any other."""
|
||||
try:
|
||||
text = path.read_text().strip()
|
||||
except OSError:
|
||||
return {}
|
||||
try:
|
||||
value = json.loads(text)
|
||||
except ValueError:
|
||||
return {'digest': text}
|
||||
return value if isinstance(value, dict) else {}
|
||||
|
||||
|
||||
def superseded(record, existing, stamp):
|
||||
"""Whether the container carries a copy this host did not write, with a
|
||||
record that differs from the one of this host. That is a container that
|
||||
came back: migrated to another node and changed there, restored from a
|
||||
backup made before the last operation, or rolled back to a snapshot. Its
|
||||
content is the one its own copy describes."""
|
||||
if not stamp.get('generation') or existing.get('generation') == stamp['generation']:
|
||||
return False
|
||||
carried_record = existing['record']
|
||||
return (carried_record.get('installation_id') == record['installation_id']
|
||||
and carried_record.get('vmid') == record['vmid'] and carried_record != record)
|
||||
|
||||
|
||||
def carry(root, vmid, mount_stopped=True, verify=False, adopted=False):
|
||||
"""Leave the current record inside its container. Returns 'carried',
|
||||
'current' when the copy was already up to date, 'skipped' when the
|
||||
instance is in the middle of an operation or cannot be reached, or
|
||||
'stale' when the container came back with another record: the record of
|
||||
this host is set aside, and the container is then recovered like any
|
||||
restored one. A stopped container is opened only when its copy is known to
|
||||
be old, or when `verify` asks to look anyway before an operation. A record
|
||||
that was just recovered is `adopted`: it replaces the copies it came from."""
|
||||
try:
|
||||
record = instances.read(root, vmid)
|
||||
except (OSError, ValueError, KeyError):
|
||||
return 'skipped'
|
||||
if record.get('status') != 'installed' or record.get('pending_transaction') \
|
||||
or record.get('pending_stack_transaction'):
|
||||
return 'skipped'
|
||||
result = _run('pct', 'config', str(vmid))
|
||||
if result.returncode != 0 or instances.identity(result.stdout.encode()) != record['installation_id']:
|
||||
return 'skipped'
|
||||
value = bundle(record, result.stdout)
|
||||
expected = digest(value)
|
||||
record_path = instances.location(root, vmid)
|
||||
path = record_path.parent / STAMP
|
||||
stamp = read_stamp(path)
|
||||
if not adopted and replaced_elsewhere(record, record_path, read_cluster(vmid), value['node']):
|
||||
record_path.replace(record_path.with_name(f"retired-{record['installation_id']}.json"))
|
||||
with contextlib.suppress(OSError):
|
||||
path.unlink()
|
||||
return 'stale'
|
||||
if not running_pid(vmid):
|
||||
if stamp.get('digest') == expected and stamp.get('generation') and not verify:
|
||||
keep_cluster(vmid, value, stamp['generation'])
|
||||
return 'current'
|
||||
if not mount_stopped:
|
||||
return 'skipped'
|
||||
try:
|
||||
with container_root(vmid, mount_stopped) as rootfs:
|
||||
if rootfs is None:
|
||||
return 'skipped'
|
||||
existing = read_copy(rootfs)
|
||||
if existing is not None and not adopted and superseded(record, existing, stamp):
|
||||
outcome = 'stale'
|
||||
elif existing is not None and digest(existing) == expected and existing.get('generation') \
|
||||
and existing.get('generation') == stamp.get('generation'):
|
||||
outcome = 'current'
|
||||
else:
|
||||
value['saved_at'] = datetime.datetime.now(datetime.timezone.utc).isoformat()
|
||||
value['generation'] = uuid.uuid4().hex
|
||||
write_copy(rootfs, value, mapped_root(result.stdout))
|
||||
stamp = {'generation': value['generation']}
|
||||
outcome = 'carried'
|
||||
if outcome == 'stale':
|
||||
record_path.replace(record_path.with_name(f"retired-{record['installation_id']}.json"))
|
||||
with contextlib.suppress(OSError):
|
||||
path.unlink()
|
||||
return outcome
|
||||
path.write_text(json.dumps({'digest': expected, 'generation': stamp.get('generation')}) + '\n')
|
||||
path.chmod(0o600)
|
||||
if stamp.get('generation'):
|
||||
keep_cluster(vmid, value, stamp['generation'])
|
||||
except OSError:
|
||||
return 'skipped'
|
||||
return outcome
|
||||
|
||||
|
||||
def registered(root):
|
||||
if not root.is_dir():
|
||||
return []
|
||||
return sorted(int(d.name) for d in root.iterdir()
|
||||
if d.name.isdecimal() and instances.has_contract(root, int(d.name)))
|
||||
|
||||
|
||||
def sync(root, vmids=None, mount_stopped=True, verify=False):
|
||||
"""Bring the copies of the given instances, or of all, up to date. Never
|
||||
waits for the registry: an operation in progress carries its own copy
|
||||
when it ends."""
|
||||
try:
|
||||
with instances.locked(root):
|
||||
return {vmid: carry(root, vmid, mount_stopped, verify)
|
||||
for vmid in (registered(root) if vmids is None else vmids)}
|
||||
except (BlockingIOError, OSError, ValueError):
|
||||
return {}
|
||||
|
||||
|
||||
def main():
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument('vmid', type=int, nargs='*')
|
||||
parser.add_argument('--root', type=Path, default=instances.ROOT)
|
||||
parser.add_argument('--running-only', action='store_true')
|
||||
parser.add_argument('--verify', action='store_true')
|
||||
args = parser.parse_args()
|
||||
if os.geteuid() != 0:
|
||||
parser.error(translate('Root privileges are required'))
|
||||
print(json.dumps(sync(args.root, args.vmid or None, not args.running_only, args.verify)))
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
sys.exit(main())
|
||||
@@ -114,9 +114,41 @@ configure_immich_ml_gpu() {
|
||||
esac
|
||||
}
|
||||
|
||||
# The generation of the AMD GPU as the compute driver names it, e.g. gfx1030.
|
||||
amd_gfx_arch() {
|
||||
local target
|
||||
target=$(awk '$1 == "gfx_target_version" && $2 > 0 {print $2}' \
|
||||
/sys/class/kfd/kfd/topology/nodes/*/properties 2>/dev/null | sort -n | tail -1)
|
||||
[[ -n $target ]] || return 1
|
||||
printf 'gfx%d%d%x' $((target / 10000)) $((target / 100 % 100)) $((target % 100))
|
||||
}
|
||||
|
||||
validate_immich_ml_runtime() {
|
||||
[[ $ML_ACCELERATION != cpu ]] || return 0
|
||||
msg_info "$(translate "Checking the GPU of the machine learning container...")"
|
||||
if [[ $ML_ACCELERATION == rocm ]]; then
|
||||
# The provider being present says nothing about this GPU: the image
|
||||
# carries kernels for some generations, and on any other every inference
|
||||
# aborts.
|
||||
local arch
|
||||
if [[ -n ${ML_GFX_OVERRIDE:-} ]]; then
|
||||
# ROCm is told to treat this GPU as the generation of its family.
|
||||
arch="gfx${ML_GFX_OVERRIDE//./}"
|
||||
else
|
||||
arch=$(amd_gfx_arch) || arch=""
|
||||
fi
|
||||
if [[ -n $arch ]]; then
|
||||
# Listed first: a match ends grep early, and with pipefail the listing
|
||||
# it cut short would read as a failure.
|
||||
local kernels
|
||||
kernels=$(pct exec "$ML_ID" -- sh -c 'ls /opt/rocm/lib/rocblas/library 2>/dev/null' || true)
|
||||
grep -Eq "[_-]${arch}\\.(dat|co|hsaco)" <<<"$kernels" || {
|
||||
oci_log "The ROCm image has no kernels for $arch"
|
||||
msg_warn "$(translate "The ROCm image has no support for the AMD GPU of this host:") $arch"
|
||||
return 1
|
||||
}
|
||||
fi
|
||||
fi
|
||||
oci_quiet pct exec "$ML_ID" -- python -c '
|
||||
import ctypes
|
||||
import sys
|
||||
@@ -135,5 +167,12 @@ else:
|
||||
assert driver.cuInit(0) == 0, "CUDA driver initialization failed"
|
||||
print("Immich ML GPU runtime:", profile, "available; model inference is tested separately")
|
||||
' "$ML_ACCELERATION" || return
|
||||
if [[ $ML_ACCELERATION == rocm ]]; then
|
||||
# One real inference on the GPU: the provider alone does not prove it.
|
||||
python3 "${SCRIPT_DIR}/oci_rocm_check.py" "$ML_ID" >>"${OCI_LOG:-/dev/null}" 2>&1 || {
|
||||
msg_warn "$(translate "The AMD GPU did not complete a test inference with ROCm")"
|
||||
return 1
|
||||
}
|
||||
fi
|
||||
msg_ok "$(translate "GPU available for machine learning:") $ML_ACCELERATION"
|
||||
}
|
||||
|
||||
@@ -0,0 +1,307 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Change what runs the recognition of an installed Immich: the CPU or a GPU.
|
||||
|
||||
The machine learning container takes the image built for the new choice and
|
||||
the devices that choice needs. Its model cache, the library, the database and
|
||||
the other containers keep their data. The change is applied with an update of
|
||||
the application, so a failure restores the previous containers, and with them
|
||||
the previous choice.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import os
|
||||
from pathlib import Path
|
||||
import re
|
||||
import subprocess
|
||||
import sys
|
||||
|
||||
import oci_instances as instances
|
||||
import oci_stack_modify
|
||||
import oci_stack_native
|
||||
from oci_ui import msg_error, msg_info, msg_ok, msg_warn, translate
|
||||
|
||||
ADAPTER = 'install_immich_stack.sh'
|
||||
ENGINE = Path(__file__).resolve().parent
|
||||
PROFILES = ('cpu', 'openvino', 'cuda', 'rocm')
|
||||
RENDER = re.compile(r'/dev/dri/renderD[0-9]+')
|
||||
# What a previous choice left in the container: its devices and the settings
|
||||
# of its runtime.
|
||||
GPU_DEVICE = re.compile(r'/dev/dri/renderD[0-9]+|/dev/kfd|/dev/nvidia\S*')
|
||||
GPU_LINE = re.compile(r'lxc\.environment\.runtime: HSA_[A-Z_]+=|lxc\.environment: NVIDIA_[A-Z_]+=|'
|
||||
r'lxc\.hook\.mount: /usr/local/lib/proxmenux/oci/nvidia-mount-')
|
||||
ROCM_ROOTFS_GB = 40
|
||||
NVIDIA_CAPABILITIES = 'compute,utility'
|
||||
GPU_VARIABLES = ('NVIDIA_DRIVER_CAPABILITIES', 'HSA_OVERRIDE_GFX_VERSION', 'HSA_USE_SVM')
|
||||
|
||||
|
||||
def run(*command):
|
||||
result = subprocess.run(command, capture_output=True, text=True, check=False)
|
||||
if result.returncode != 0:
|
||||
detail = (result.stderr or result.stdout).strip().splitlines()[-1:] or ['']
|
||||
raise RuntimeError(f"{' '.join(command[:3])}: {detail[0]}")
|
||||
return result.stdout
|
||||
|
||||
|
||||
def members(root, primary_id):
|
||||
"""The record of each container of the application, by its role."""
|
||||
primary = instances.read(root, primary_id)
|
||||
found = {}
|
||||
for member in (primary.get('stack') or {}).get('members', []):
|
||||
record = instances.read(root, member['vmid'])
|
||||
profile = record['deployment'].get('replay_profile') or {}
|
||||
if profile.get('adapter') != ADAPTER:
|
||||
raise ValueError(translate('This application is not an Immich installed by ProxMenux'))
|
||||
found[profile.get('role')] = record
|
||||
if set(found) != {'server', 'database', 'valkey', 'machine-learning'} or found['server']['vmid'] != primary_id:
|
||||
raise ValueError(translate('This application is not an Immich installed by ProxMenux'))
|
||||
return primary, found
|
||||
|
||||
|
||||
def current(record):
|
||||
return (record['deployment'].get('machine_learning') or {}).get('acceleration', 'cpu')
|
||||
|
||||
|
||||
def conf(vmid):
|
||||
return Path(f'/etc/pve/lxc/{int(vmid)}.conf')
|
||||
|
||||
|
||||
def settings(vmid):
|
||||
"""The current configuration of the container, one value per key."""
|
||||
lines = [line.partition(': ') for line in run('pct', 'config', str(vmid)).splitlines()]
|
||||
return {key: value for key, separator, value in lines if separator}
|
||||
|
||||
|
||||
def strip_gpu(vmid):
|
||||
"""Take the devices and runtime settings of the previous choice away."""
|
||||
for key, value in oci_stack_modify.entries(vmid, 'dev').items():
|
||||
if GPU_DEVICE.fullmatch(oci_stack_modify.device_path(value)):
|
||||
run('pct', 'set', str(vmid), '--delete', key)
|
||||
text = conf(vmid).read_text()
|
||||
head, separator, snapshots = text.partition('\n[')
|
||||
rebuilt = '\n'.join(line for line in head.split('\n') if not GPU_LINE.match(line)) + separator + snapshots
|
||||
if rebuilt != text:
|
||||
conf(vmid).write_text(rebuilt)
|
||||
|
||||
|
||||
def append_lines(vmid, lines):
|
||||
"""Add lines to the current configuration, before any snapshot section."""
|
||||
head, separator, snapshots = conf(vmid).read_text().partition('\n[')
|
||||
conf(vmid).write_text(head.rstrip('\n') + '\n' + '\n'.join(lines) + '\n' + separator + snapshots)
|
||||
|
||||
|
||||
def add_device(vmid, path):
|
||||
info = os.stat(path)
|
||||
run('pct', 'set', str(vmid), '--' + oci_stack_modify.free_key(vmid, 'dev'),
|
||||
f'path={path},gid={info.st_gid},mode=0660')
|
||||
|
||||
|
||||
def configure(vmid, acceleration, render_device, gfx_override):
|
||||
"""Give the stopped container what the new choice needs."""
|
||||
strip_gpu(vmid)
|
||||
if acceleration in ('openvino', 'rocm'):
|
||||
add_device(vmid, render_device)
|
||||
if acceleration == 'rocm':
|
||||
add_device(vmid, '/dev/kfd')
|
||||
if gfx_override:
|
||||
append_lines(vmid, [f'lxc.environment.runtime: HSA_OVERRIDE_GFX_VERSION={gfx_override}',
|
||||
'lxc.environment.runtime: HSA_USE_SVM=0'])
|
||||
size = re.search(r'size=([0-9]+)G', settings(vmid).get('rootfs', ''))
|
||||
if size and int(size[1]) < ROCM_ROOTFS_GB:
|
||||
# The ROCm image is several times larger than the others.
|
||||
run('pct', 'resize', str(vmid), 'rootfs', f'{ROCM_ROOTFS_GB}G')
|
||||
if acceleration == 'cuda':
|
||||
# The same NVIDIA setup the installation uses.
|
||||
run('bash', '-c', f'set -e; SCRIPT_DIR="{ENGINE}"; source "$SCRIPT_DIR/oci_ui.sh"; '
|
||||
'die() { msg_error "$*"; exit 1; }; oci_log() { :; }; '
|
||||
'source "$SCRIPT_DIR/oci_nvidia_setup.sh"; source "$SCRIPT_DIR/oci_immich_ml.sh"; '
|
||||
f'configure_immich_nvidia {int(vmid)} {NVIDIA_CAPABILITIES}')
|
||||
run('pct', 'set', str(vmid), '--memory', str(memory(acceleration)))
|
||||
values = settings(vmid)
|
||||
limited = 'cpulimit' in values and 'cores' not in values
|
||||
# Intel keeps the CPU topology and limits its share of time; the others
|
||||
# are given four cores.
|
||||
if acceleration == 'openvino' and not limited:
|
||||
run('pct', 'set', str(vmid), '--cpulimit', '4', '--delete', 'cores')
|
||||
elif acceleration != 'openvino' and limited:
|
||||
run('pct', 'set', str(vmid), '--cores', '4', '--delete', 'cpulimit')
|
||||
|
||||
|
||||
def validate(acceleration, render_device, gfx_override):
|
||||
if acceleration not in PROFILES:
|
||||
raise ValueError(translate('Immich GPU profile not validated'))
|
||||
if acceleration in ('openvino', 'rocm'):
|
||||
if not render_device or not RENDER.fullmatch(render_device) or not Path(render_device).is_char_device():
|
||||
raise ValueError(f"{translate('The render device does not exist:')} {render_device}")
|
||||
vendor = Path(f'/sys/class/drm/{Path(render_device).name}/device/vendor').read_text().strip()
|
||||
if vendor != ('0x8086' if acceleration == 'openvino' else '0x1002'):
|
||||
raise ValueError(translate('The render device does not belong to the GPU of that choice'))
|
||||
if acceleration == 'rocm' and not Path('/dev/kfd').is_char_device():
|
||||
raise ValueError(translate('ROCm requires /dev/kfd on the host'))
|
||||
if acceleration == 'cuda' and subprocess.run(['nvidia-smi', '-L'], capture_output=True, check=False).returncode:
|
||||
raise ValueError(translate('The host NVIDIA driver is not responding correctly'))
|
||||
if gfx_override and (acceleration != 'rocm' or not re.fullmatch(r'[0-9]{1,2}\.[0-9]\.[0-9]', gfx_override)):
|
||||
raise ValueError(translate('Invalid ROCm generation override'))
|
||||
|
||||
|
||||
def recorded(record, acceleration, render_device, gfx_override):
|
||||
"""The recognition settings of a record after the change."""
|
||||
settings = dict(record['deployment'].get('machine_learning') or {})
|
||||
settings.update(acceleration=acceleration,
|
||||
render_device=render_device if acceleration in ('openvino', 'rocm') else None,
|
||||
gfx_override=gfx_override if acceleration == 'rocm' else None)
|
||||
settings['resources'] = dict(settings.get('resources') or {}, memory_mb=memory(acceleration),
|
||||
cpu_allocation='quota' if acceleration == 'openvino' else 'cpuset')
|
||||
return settings
|
||||
|
||||
|
||||
def memory(acceleration):
|
||||
return 4096 if acceleration == 'cpu' else 8192
|
||||
|
||||
|
||||
def resources(current, acceleration):
|
||||
"""The resources a record declares for the container of the new choice."""
|
||||
result = dict(current, memory_mb=memory(acceleration))
|
||||
result.pop('cpu_allocation', None)
|
||||
if acceleration == 'openvino':
|
||||
result['cpu_allocation'] = 'quota'
|
||||
return result
|
||||
|
||||
|
||||
def declared(environment, acceleration, gfx_override):
|
||||
"""The variables a record declares, with the ones of the new choice."""
|
||||
result = [entry for entry in environment if entry.get('name') not in GPU_VARIABLES]
|
||||
if acceleration == 'cuda':
|
||||
result.append({'name': 'NVIDIA_DRIVER_CAPABILITIES', 'value': NVIDIA_CAPABILITIES})
|
||||
if acceleration == 'rocm' and gfx_override:
|
||||
result += [{'name': 'HSA_OVERRIDE_GFX_VERSION', 'value': gfx_override}, {'name': 'HSA_USE_SVM', 'value': '0'}]
|
||||
return result
|
||||
|
||||
|
||||
def declare(record, acceleration, gfx_override):
|
||||
"""A record an update already rebuilt describes its container by itself:
|
||||
it takes the variables, the resources and the disk of the new choice."""
|
||||
plan = record['deployment']
|
||||
if 'native_config' in plan:
|
||||
return
|
||||
plan['environment'] = declared(plan.get('environment', []), acceleration, gfx_override)
|
||||
plan['resources'] = resources(plan.get('resources') or {}, acceleration)
|
||||
size = re.search(r'size=([0-9]+)G', settings(record['vmid']).get('rootfs', ''))
|
||||
if size and isinstance(plan.get('rootfs'), dict):
|
||||
plan['rootfs']['size_gb'] = int(size[1])
|
||||
profile = (record['template'].get('proxmox') or {}).get('installer_profile')
|
||||
if isinstance(profile, dict):
|
||||
profile.pop('cpu_allocation', None)
|
||||
if acceleration == 'openvino':
|
||||
profile['cpu_allocation'] = 'quota'
|
||||
|
||||
|
||||
def chosen(record):
|
||||
"""What runs recognition according to a record."""
|
||||
saved = record['deployment'].get('machine_learning') or {}
|
||||
return saved.get('acceleration', 'cpu'), saved.get('render_device'), saved.get('gfx_override')
|
||||
|
||||
|
||||
def stop(vmid):
|
||||
try:
|
||||
run('pct', 'shutdown', str(vmid), '--timeout', '60')
|
||||
except RuntimeError:
|
||||
run('pct', 'stop', str(vmid))
|
||||
|
||||
|
||||
def apply(root, primary_id, vmid, images, acceleration, render_device, gfx_override):
|
||||
"""Give the stopped container and the records one choice."""
|
||||
configure(vmid, acceleration, render_device, gfx_override)
|
||||
learning = instances.read(root, vmid)
|
||||
learning['template']['container_contract']['image']['reference'] = images[acceleration]
|
||||
learning['deployment']['machine_learning'] = recorded(learning, acceleration, render_device, gfx_override)
|
||||
declare(learning, acceleration, gfx_override)
|
||||
instances.write(instances.location(root, vmid), learning)
|
||||
primary = instances.read(root, primary_id)
|
||||
primary['stack']['deployment']['machine_learning'] = recorded(
|
||||
{'deployment': primary['stack']['deployment']}, acceleration, render_device, gfx_override)
|
||||
intent = (primary.get('native_stack_intent') or {}).get('deployment')
|
||||
if isinstance(intent, dict) and 'machine_learning' in intent:
|
||||
intent['machine_learning'] = recorded({'deployment': intent}, acceleration, render_device, gfx_override)
|
||||
instances.write(instances.location(root, primary_id), primary)
|
||||
oci_stack_modify.register(root, vmid)
|
||||
|
||||
|
||||
def revert(root, primary_id, vmid, images, previous, was_running):
|
||||
"""Give the container and the records the choice they had."""
|
||||
try:
|
||||
if oci_stack_modify.is_running(vmid):
|
||||
stop(vmid)
|
||||
apply(root, primary_id, vmid, images, *previous)
|
||||
except (OSError, ValueError, KeyError, RuntimeError, subprocess.SubprocessError) as error:
|
||||
msg_warn(f"{translate('The previous recognition choice could not be put back:')} {error}")
|
||||
return
|
||||
if was_running:
|
||||
subprocess.run(['pct', 'start', str(vmid)], capture_output=True, check=False)
|
||||
|
||||
|
||||
def change(root, primary_id, acceleration, render_device=None, gfx_override=None):
|
||||
validate(acceleration, render_device, gfx_override)
|
||||
with instances.locked(root):
|
||||
primary, found = members(root, primary_id)
|
||||
learning = found['machine-learning']
|
||||
vmid = learning['vmid']
|
||||
if any(record.get('status') != 'installed' or record.get('pending_transaction')
|
||||
or record.get('pending_stack_transaction') for record in found.values()):
|
||||
raise ValueError(translate('The container has an operation pending; finish or recover it first'))
|
||||
if current(learning) == acceleration:
|
||||
raise ValueError(translate('Recognition already runs on that choice'))
|
||||
images = primary['stack']['template']['proxmox']['application_options']['machine_learning']['profile_images']
|
||||
previous = chosen(learning)
|
||||
was_running = oci_stack_modify.is_running(vmid)
|
||||
try:
|
||||
if was_running:
|
||||
msg_info(translate('Stopping the container...'))
|
||||
stop(vmid)
|
||||
msg_ok(translate('Container stopped'))
|
||||
msg_info(translate('Preparing the machine learning container for the new choice...'))
|
||||
apply(root, primary_id, vmid, images, acceleration, render_device, gfx_override)
|
||||
msg_ok(translate('Machine learning container prepared'))
|
||||
except BaseException:
|
||||
revert(root, primary_id, vmid, images, previous, was_running)
|
||||
raise
|
||||
try:
|
||||
# The image of the new choice replaces the one in use, as an update does.
|
||||
oci_stack_native.run(primary_id, acknowledge_external_data=True)
|
||||
except BaseException:
|
||||
with instances.locked(root):
|
||||
# An update left halfway keeps its own record of what to restore.
|
||||
if not instances.read(root, primary_id).get('pending_stack_transaction'):
|
||||
revert(root, primary_id, vmid, images, previous, was_running)
|
||||
raise
|
||||
if was_running and not oci_stack_modify.is_running(vmid):
|
||||
# An update leaves each container as it found it, and this one was
|
||||
# stopped for the change.
|
||||
run('pct', 'start', str(vmid))
|
||||
|
||||
|
||||
def main():
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument('vmid', type=int, help='main container of the application')
|
||||
parser.add_argument('--acceleration', required=True, choices=PROFILES)
|
||||
parser.add_argument('--render-device')
|
||||
parser.add_argument('--gfx-override')
|
||||
parser.add_argument('--root', type=Path, default=instances.ROOT)
|
||||
args = parser.parse_args()
|
||||
if os.geteuid() != 0:
|
||||
parser.error(translate('Root privileges are required'))
|
||||
try:
|
||||
change(args.root, args.vmid, args.acceleration, args.render_device, args.gfx_override)
|
||||
except BlockingIOError:
|
||||
msg_error(translate('Another OCI operation is using the instance registry. Wait for it to finish.'))
|
||||
return 1
|
||||
except (OSError, ValueError, KeyError, RuntimeError, StopIteration, subprocess.SubprocessError) as error:
|
||||
if not getattr(error, 'oci_reported', False):
|
||||
msg_error(f"{translate('The recognition of Immich was not changed:')} {error}")
|
||||
return 1
|
||||
msg_ok(translate('Recognition of Immich changed; the application was updated with the new choice.'))
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
sys.exit(main())
|
||||
@@ -88,6 +88,8 @@ def image_from_archive(path):
|
||||
manifest = blob(digest)
|
||||
config = blob(manifest['config']['digest'])
|
||||
return {'manifest_digest': digest, 'config_digest': manifest['config']['digest'],
|
||||
'layers': [layer['digest'] for layer in manifest.get('layers', [])],
|
||||
'created': config.get('created'),
|
||||
'architecture': config['architecture'], 'os': config.get('os'),
|
||||
'defaults': {k: config.get('config', {}).get(k) for k in FIELDS}}
|
||||
|
||||
@@ -190,8 +192,20 @@ def resolve_candidate(reference, architecture):
|
||||
# The build date identifies the image as the publisher released it: it is
|
||||
# what changes when an image is rebuilt, whether or not the application
|
||||
# version inside it moved.
|
||||
return {'manifest_digest': digest, 'defaults': defaults, 'version': version,
|
||||
'created': config.get('created')}
|
||||
return {'manifest_digest': digest, 'layers': [layer['digest'] for layer in json.loads(raw).get('layers', [])],
|
||||
'defaults': defaults, 'version': version, 'created': config.get('created')}
|
||||
|
||||
|
||||
def same_image(candidate, image):
|
||||
"""Whether what the registry serves is the image of an archive or a record.
|
||||
|
||||
The OCI archive made from a manifest published in the Docker format has
|
||||
another manifest and another configuration digest: both are rewritten. Its
|
||||
layers and the moment the image was built are the same in both."""
|
||||
if candidate['manifest_digest'] == image.get('manifest_digest'):
|
||||
return True
|
||||
return bool(candidate.get('layers') and candidate.get('created')
|
||||
and candidate['layers'] == image.get('layers') and candidate['created'] == image.get('created'))
|
||||
|
||||
|
||||
def compare(record, current, candidate=None):
|
||||
|
||||
@@ -764,9 +764,11 @@ def restore_description(vmid, state):
|
||||
run('pct', 'set', str(vmid), '--description', original_description(state))
|
||||
|
||||
|
||||
def release_stage(state):
|
||||
def release_stage(state, discard=False):
|
||||
"""After a commit the holder CT only keeps its own rootfs: every parked
|
||||
volume went back to the application. Anything still attached keeps it."""
|
||||
volume went back to the application. After a verified recovery it holds
|
||||
the disks of the failed attempt, which the restored backup replaced:
|
||||
`discard` removes them with it. Otherwise anything still attached keeps it."""
|
||||
stage = state.get('stage')
|
||||
if not stage:
|
||||
return
|
||||
@@ -774,11 +776,23 @@ def release_stage(state):
|
||||
config = owned(stage, 'proxmenux-transaction=' + state['id'])
|
||||
except (ValueError, RuntimeError, subprocess.CalledProcessError):
|
||||
return
|
||||
if mounts(config) or any(re.fullmatch(r'unused[0-9]+', key) for key in parse_config(config)):
|
||||
held = mounts(config) or any(re.fullmatch(r'unused[0-9]+', key) for key in parse_config(config))
|
||||
if held and not discard:
|
||||
return
|
||||
if discard:
|
||||
stop(stage)
|
||||
run('pct', 'destroy', str(stage))
|
||||
|
||||
|
||||
def discard_stage(state):
|
||||
"""Remove the holder CT of an operation that ended in a verified recovery.
|
||||
A failure here leaves it in place and does not undo the recovery."""
|
||||
try:
|
||||
release_stage(state, discard=True)
|
||||
except (OSError, ValueError, RuntimeError, subprocess.SubprocessError) as error:
|
||||
log(f'cleanup: {error}')
|
||||
|
||||
|
||||
def gib(size):
|
||||
return f'{size / 1024**3:.1f} GB'
|
||||
|
||||
@@ -1124,7 +1138,9 @@ def complete_recovery(root, journal, state):
|
||||
instances.write(instances.location(root, vmid), restored)
|
||||
checkpoint(journal, state, 'rolled-back')
|
||||
if show:
|
||||
msg_ok(translate('Recovery completed. The displaced disks and the backup are kept; nothing was deleted automatically.'))
|
||||
# A member of a stack is released when the whole stack is back.
|
||||
discard_stage(state)
|
||||
msg_ok(translate('Recovery completed. The disks of the failed attempt were removed; the backup is kept.'))
|
||||
if state.get('original_host_sources') or state.get('desired_host_sources'):
|
||||
msg_info2(translate('Shared host files are kept as they are; the backup does not restore their content.'))
|
||||
|
||||
|
||||
@@ -2,12 +2,14 @@
|
||||
"""Instance recording adapter for the existing dedicated stack installers."""
|
||||
import argparse
|
||||
import copy
|
||||
import ipaddress
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
import re
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
|
||||
import oci_console
|
||||
import oci_instances as instances
|
||||
@@ -103,6 +105,40 @@ def capture_rootfs(root, vmid):
|
||||
instances.write(instances.location(root, vmid), record)
|
||||
|
||||
|
||||
PRIVATE_STACK_NETWORK = ipaddress.ip_network('10.77.0.0/16')
|
||||
|
||||
|
||||
def access_address(vmid, wait=20):
|
||||
"""The address a container is reached at from the local network. A member
|
||||
of a multi-container application also has a leg on its private network,
|
||||
and that address answers from the host only: it is used when the
|
||||
container has no other."""
|
||||
def addresses():
|
||||
result = subprocess.run(['lxc-info', '-n', str(vmid), '-iH'], capture_output=True, text=True, timeout=5)
|
||||
found = []
|
||||
for line in result.stdout.splitlines():
|
||||
try:
|
||||
found.append(ipaddress.IPv4Address(line.strip()))
|
||||
except ValueError:
|
||||
continue
|
||||
return found
|
||||
|
||||
def outside(found):
|
||||
return next((str(address) for address in found if address not in PRIVATE_STACK_NETWORK), '')
|
||||
|
||||
found = addresses()
|
||||
config = subprocess.run(['pct', 'config', str(vmid)], capture_output=True, text=True, timeout=30).stdout
|
||||
legs = re.findall(r'^net[0-9]+: .*?(?:^|,)ip=([^,\s]+)', config, re.MULTILINE)
|
||||
# A leg outside the private network may still be waiting for its lease.
|
||||
expected = any(leg == 'dhcp' or (leg[:1].isdigit() and ipaddress.ip_interface(leg).ip not in PRIVATE_STACK_NETWORK)
|
||||
for leg in legs)
|
||||
deadline = time.monotonic() + (wait if expected else 0)
|
||||
while not outside(found) and time.monotonic() < deadline:
|
||||
time.sleep(2)
|
||||
found = addresses()
|
||||
return outside(found) or (str(found[0]) if found else '')
|
||||
|
||||
|
||||
def finalize(root, primary):
|
||||
parent = instances.read(root, primary)
|
||||
intent = parent['native_stack_intent']
|
||||
@@ -118,10 +154,7 @@ def finalize(root, primary):
|
||||
presentation['catalog_ui'] = {**stack_ui, **(presentation.get('catalog_ui') or {})}
|
||||
if not (presentation.get('first_run') or {}).get('endpoints'):
|
||||
presentation['first_run'] = intent['template'].get('first_run') or {}
|
||||
ip_result = subprocess.run(['lxc-info', '-n', str(vmid), '-iH'],
|
||||
capture_output=True, text=True, timeout=5)
|
||||
ip = next((line.strip() for line in ip_result.stdout.splitlines()
|
||||
if re.fullmatch(r'[0-9]+(?:\.[0-9]+){3}', line.strip())), '')
|
||||
ip = access_address(vmid)
|
||||
description = render(presentation, plan['image']['manifest_digest'],
|
||||
record['installation_id'], ip)
|
||||
subprocess.run(['pct', 'set', str(vmid), '--description', description], check=True)
|
||||
|
||||
@@ -68,6 +68,36 @@ def create_parent(path, root):
|
||||
os.chown(path, owner.st_uid, owner.st_gid)
|
||||
|
||||
|
||||
def rebuild(root, vmid, same_gpu=False):
|
||||
"""Point a stopped container at the NVIDIA driver of this host. The caller
|
||||
holds the registry and the container is not started: it is the step a
|
||||
restore on a host with another driver needs before the first start.
|
||||
Returns whether anything had to change."""
|
||||
record = instances.read(root, vmid)
|
||||
config = instances.command('pct', 'config', str(vmid))
|
||||
if not instances.same_config_except_notes(record, config):
|
||||
raise ValueError(translate('The container identity or configuration changed'))
|
||||
plan = nv.refresh_plan(config, record['observed']['gpu_devices'][nv.KEY], same_gpu=same_gpu)
|
||||
if not plan['changed']:
|
||||
return False
|
||||
if instances.command('pct', 'status', str(vmid)).strip() != b'status: stopped':
|
||||
raise ValueError(translate('Stop the container before the NVIDIA refresh'))
|
||||
instances.command('pct', 'mount', str(vmid))
|
||||
try:
|
||||
candidate = prepare(Path(f'/var/lib/lxc/{vmid}/rootfs'), plan)
|
||||
finally:
|
||||
instances.command('pct', 'unmount', str(vmid))
|
||||
Path(f'/etc/pve/lxc/{vmid}.conf').write_bytes(candidate)
|
||||
updated = copy.deepcopy(record)
|
||||
updated['observed'] = instances.observe(vmid, record['installation_id'],
|
||||
record['observed']['archive_path'], record['observed']['resolved_registry_digest'],
|
||||
record['observed']['image'])
|
||||
nv.check_mounts(updated['observed']['config'].encode(), plan['inventory'])
|
||||
nv.check_devices(updated['observed']['config'].encode(), plan['inventory'])
|
||||
instances.write(instances.location(root, vmid), updated)
|
||||
return True
|
||||
|
||||
|
||||
def refresh(root, vmid, apply=False):
|
||||
with instances.locked(root):
|
||||
record = instances.read(root, vmid)
|
||||
|
||||
@@ -75,11 +75,13 @@ def verify(value):
|
||||
raise ValueError(translate('The NVIDIA driver or inventory changed; the operation was stopped'))
|
||||
|
||||
|
||||
def refresh_plan(config, previous, current=None):
|
||||
def refresh_plan(config, previous, current=None, same_gpu=True):
|
||||
"""Resolve current host components without treating a driver version as intent.
|
||||
|
||||
This only prepares a plan; applying it requires a stopped-CT transaction and
|
||||
preparing file destinations/library links before the next native start.
|
||||
A container restored on another host takes the GPU of that host: the
|
||||
application asked for the NVIDIA runtime, not for one card.
|
||||
"""
|
||||
check_devices(config, previous)
|
||||
check_mounts(config, previous)
|
||||
@@ -92,7 +94,7 @@ def refresh_plan(config, previous, current=None):
|
||||
raise ValueError(translate('Incomplete NVIDIA identity'))
|
||||
result.append(tuple(fields[:2]))
|
||||
return sorted(result)
|
||||
if identities(previous) != identities(current):
|
||||
if same_gpu and identities(previous) != identities(current):
|
||||
raise ValueError(translate('The physical NVIDIA selection changed'))
|
||||
# Remove only entries already validated against our recorded inventory.
|
||||
kept = []
|
||||
|
||||
@@ -29,6 +29,7 @@ CLUSTER_NODES = Path('/etc/pve/nodes')
|
||||
SNIPPETS = Path('/var/lib/vz/snippets')
|
||||
# The App tab of ProxMenux Monitor keeps one file per VMID.
|
||||
MONITOR_APPS = Path('/etc/proxmenux/apps')
|
||||
CLUSTER_RECORDS = Path('/etc/pve/priv/proxmenux/oci')
|
||||
HOST_MONITOR_INCLUDES = (Path('/etc/pve/proxmenux/host-monitor'), Path('/etc/pve/lxc/proxmenux-host-monitor'))
|
||||
STACK_HOOK = 'proxmenux-stack-dependencies.sh'
|
||||
|
||||
@@ -136,6 +137,7 @@ def remove_host_state(vmid):
|
||||
except OSError:
|
||||
pass
|
||||
_unlink(hook)
|
||||
_unlink(CLUSTER_RECORDS / f'{int(vmid)}.json')
|
||||
_unlink(MONITOR_APPS / f'{int(vmid)}.json')
|
||||
dismissed = MONITOR_APPS / '.oci-dismissed.json'
|
||||
try:
|
||||
@@ -178,6 +180,7 @@ def _leftovers(vmid):
|
||||
"""Whether anything of the container is still on the host."""
|
||||
paths = [runtime_settings.include_path(vmid), runtime_settings.legacy_include_path(vmid),
|
||||
SNIPPETS / f'proxmenux-rclone-{int(vmid)}-fuse-hook.sh', MONITOR_APPS / f'{int(vmid)}.json',
|
||||
CLUSTER_RECORDS / f'{int(vmid)}.json',
|
||||
*oci_console.LOG_DIR.glob(f'{int(vmid)}.console.log*')]
|
||||
return any(path.exists() for path in paths)
|
||||
|
||||
@@ -191,7 +194,8 @@ def sweep_orphans(root):
|
||||
for directory, pattern in ((runtime_settings.include_path(0).parent, r'([0-9]+)\.sysctls'),
|
||||
(runtime_settings.legacy_include_path(0).parent, r'([0-9]+)\.proxmenux-sysctls'),
|
||||
(oci_console.LOG_DIR, r'([0-9]+)\.console\.log.*'),
|
||||
(SNIPPETS, r'proxmenux-rclone-([0-9]+)-fuse-hook\.sh')):
|
||||
(SNIPPETS, r'proxmenux-rclone-([0-9]+)-fuse-hook\.sh'),
|
||||
(CLUSTER_RECORDS, r'([0-9]+)\.json')):
|
||||
if directory.is_dir():
|
||||
found.update(int(m.group(1)) for m in (re.fullmatch(pattern, p.name) for p in directory.iterdir()) if m)
|
||||
# Only the App tab registrations of OCI installs; the other ones belong to
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,81 @@
|
||||
#!/usr/bin/env python3
|
||||
"""A real inference on the AMD GPU of a container, through ROCm.
|
||||
|
||||
That the MIGraphX provider is listed says nothing about the GPU of this host:
|
||||
an image carries kernels for some generations only, and on any other the
|
||||
first inference aborts. This runs a small convolution model on the GPU inside
|
||||
the container and fails when it does not come back with a result.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import re
|
||||
import subprocess
|
||||
import sys
|
||||
|
||||
# A model of five operators (convolution, pooling and a dense layer), enough
|
||||
# to make ROCm compile and run on the GPU without downloading anything.
|
||||
MODEL = (
|
||||
'CAg65wUKPAoBeAoBdwoBYhIBYyIEQ29udioVCgxrZXJuZWxfc2hhcGVAA0ADoAEHKhEKBHBhZHNAAUABQAFAAaABBwoMCgFjEgFy'
|
||||
'IgRSZWx1ChkKAXISAWciEUdsb2JhbEF2ZXJhZ2VQb29sCg8KAWcSAWYiB0ZsYXR0ZW4KFAoBZgoCZncKAmZiEgF5IgRHZW1tEgVw'
|
||||
'cm9iZSrAAwgECAMIAwgDEAFCAXdKsAOu/QA5e7v0POCS4LyqZLa9sDs6vdcWy70cFMU78DwJPpibSb2CJX69qqNIPVEuEj3xtSw8'
|
||||
'U4++vWu0P7vqZY49x6UJvn5wO71qr0K+dgwEvvuXPL4vlsC8WskBvkI43jwWaYA8QyKZvKzbgL4Lply9i+2euzylOTyXrxy+ELBD'
|
||||
'vZVmyL1epqW9pUXZPRNipb1fIlW7gB+1PfIKb70yAze8Bfw0PAgA0Ts15Pq9D3/5O74kCz54bR6+ZwCwPbWMQzyGX4O9uNdMPl0c'
|
||||
'nD1HnfW9vyz0O0o2bD17ppq8K9yLPcf22bv9pog9Ak4TPipgir1BaaY8U8U9vT+EUDwwI/O9LUhtvUe5oLwdEbg9nYrqPX2HB74m'
|
||||
'vqK9X3yEPRcGTL7itj29GGUfvOW3AD6fMI090AYGvfz3Fr3I9cy8aQIcPqtRL71lxvi8pWsQPczeRbyAnaG8NSnkvaIDl7rdsDW9'
|
||||
'r9LuPabAhT1DOh67a+KIPeg1C726edc914sNuhP0bj3+LwS+CgAOPUPfLL7talC+b235vCBOuL1eZIY889xlPj9Wqr06kX+9VESo'
|
||||
'PDHwST0qGQgEEAFCAWJKEAAAAAAAAAAAAAAAAAAAAAAqLAgECAIQAUICZndKIAmDkLy5sqi8S92PPUX0VD1istO9DrsBvIJBZztd'
|
||||
'9de9KhIIAhABQgJmYkoIAAAAAAAAAABaGwoBeBIWChQIARIQCgIIAQoCCAMKAggICgIICGITCgF5Eg4KDAgBEggKAggBCgIIAkIE'
|
||||
'CgAQDQ=='
|
||||
)
|
||||
PROBE = """
|
||||
import base64, sys
|
||||
import numpy as np
|
||||
import onnxruntime as ort
|
||||
options = ort.SessionOptions()
|
||||
options.log_severity_level = 3
|
||||
session = ort.InferenceSession(base64.b64decode(sys.argv[1]), options, providers=['MIGraphXExecutionProvider'])
|
||||
assert session.get_providers()[0] == 'MIGraphXExecutionProvider', session.get_providers()
|
||||
result = session.run(None, {'x': np.ones((1, 3, 8, 8), np.float32)})[0]
|
||||
assert result.shape == (1, 2) and np.isfinite(result).all(), result
|
||||
print('ROCm inference on the GPU: ok')
|
||||
"""
|
||||
|
||||
|
||||
def variables(config):
|
||||
"""What the container tells ROCm about its GPU. A command run in the
|
||||
container does not receive the variables its application starts with."""
|
||||
found = re.findall(r'^lxc\.environment\.runtime: (.+)$', config, re.M)
|
||||
for line in re.findall(r'^env: (.+)$', config, re.M):
|
||||
found += line.split('\0')
|
||||
return [value for value in found if re.fullmatch(r'HSA_[A-Z_]+=[0-9A-Za-z.]+', value)]
|
||||
|
||||
|
||||
def check(vmid, timeout=300):
|
||||
"""Raises RuntimeError when the container cannot run the model on its GPU."""
|
||||
vmid = str(int(vmid))
|
||||
config = subprocess.run(['pct', 'config', vmid], capture_output=True, text=True, check=False).stdout
|
||||
try:
|
||||
result = subprocess.run(['pct', 'exec', vmid, '--', 'env', *variables(config), 'python', '-c', PROBE, MODEL],
|
||||
capture_output=True, text=True, timeout=timeout, check=False)
|
||||
except subprocess.TimeoutExpired as error:
|
||||
raise RuntimeError('the GPU did not answer in time') from error
|
||||
if result.returncode != 0:
|
||||
detail = (result.stderr or result.stdout).strip().splitlines()[-1:] or ['']
|
||||
raise RuntimeError(detail[0][:300])
|
||||
|
||||
|
||||
def main():
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument('vmid', type=int)
|
||||
args = parser.parse_args()
|
||||
try:
|
||||
check(args.vmid)
|
||||
except RuntimeError as error:
|
||||
print(error, file=sys.stderr)
|
||||
return 1
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
sys.exit(main())
|
||||
@@ -377,6 +377,8 @@ class NativeAdapter:
|
||||
member_tx.run('pct', 'exec', str(vmid), '--', 'python', '-c',
|
||||
'import onnxruntime as ort; '
|
||||
'assert "MIGraphXExecutionProvider" in ort.get_available_providers()')
|
||||
import oci_rocm_check
|
||||
oci_rocm_check.check(vmid)
|
||||
if acceleration in ('openvino', 'cuda'):
|
||||
member_tx.run('pct', 'exec', str(vmid), '--', 'python', '-c',
|
||||
'import sys,ctypes,onnxruntime as ort; p=sys.argv[1]; '
|
||||
@@ -561,7 +563,7 @@ class NativeAdapter:
|
||||
record.pop('pending_stack_transaction', None)
|
||||
instances.write(instances.location(self.root, vmid), record)
|
||||
try:
|
||||
self.release_stages()
|
||||
self.release_stages(recovered=state['phase'] == 'rolled-back')
|
||||
if self.keep_backup and state['phase'] == 'committed':
|
||||
import oci_keep_backup
|
||||
for backup in sorted(self.journal.parent.glob('backup-*/vzdump-lxc-*.tar.zst')):
|
||||
@@ -574,18 +576,25 @@ class NativeAdapter:
|
||||
except (OSError, ValueError) as exc:
|
||||
member_tx.log(f'cleanup: {exc}')
|
||||
|
||||
def release_stages(self):
|
||||
"""The temporary containers that held the data of each member; one
|
||||
that still has a disk attached is kept."""
|
||||
def release_stages(self, recovered=False):
|
||||
"""The temporary containers that held the data of each member. When
|
||||
the whole stack was recovered they are removed with the disks of the
|
||||
failed attempt; otherwise one that still has a disk attached is kept."""
|
||||
journals = [(path, True) for path in self.journal.parent.glob('recovery-*/transaction.json')]
|
||||
for vmid in self.records:
|
||||
folder = instances.location(self.root, vmid).parent / 'transactions'
|
||||
for member_journal in folder.glob('*/transaction.json'):
|
||||
try:
|
||||
state = json.loads(member_journal.read_text())
|
||||
except (OSError, ValueError):
|
||||
continue
|
||||
if (state.get('coordinated') or {}).get('journal') == str(self.journal):
|
||||
member_tx.release_stage(state)
|
||||
journals += [(path, False) for path in folder.glob('*/transaction.json')]
|
||||
for member_journal, own in journals:
|
||||
try:
|
||||
state = json.loads(member_journal.read_text())
|
||||
except (OSError, ValueError):
|
||||
continue
|
||||
if not own and (state.get('coordinated') or {}).get('journal') != str(self.journal):
|
||||
continue
|
||||
if recovered and state.get('phase') == 'rolled-back':
|
||||
member_tx.discard_stage(state)
|
||||
else:
|
||||
member_tx.release_stage(state)
|
||||
|
||||
def prune_backups(self, include_current):
|
||||
"""The backups of closed operations are removed, those of this one
|
||||
|
||||
@@ -12,7 +12,7 @@ import tempfile
|
||||
|
||||
import oci_instances as instances
|
||||
import oci_instance_transaction as transaction
|
||||
from oci_installation_state import image_from_archive
|
||||
from oci_installation_state import image_from_archive, same_image
|
||||
from oci_ui import translate, msg_info, msg_ok, msg_warn, msg_error, msg_info2
|
||||
|
||||
|
||||
@@ -65,7 +65,7 @@ def resolve_archive(desired, config, current=None, check=None):
|
||||
if not re.fullmatch(r'sha256:[a-f0-9]{64}', digest):
|
||||
raise ValueError(translate('Invalid registry digest'))
|
||||
msg_ok(f"{translate('Image:')} {reference} ({candidate.get('version') or digest[7:19]})")
|
||||
if digest == current:
|
||||
if current and same_image(candidate, current):
|
||||
return None, digest
|
||||
if check:
|
||||
check()
|
||||
@@ -85,7 +85,7 @@ def resolve_archive(desired, config, current=None, check=None):
|
||||
translate('The image did not pass the integrity check'))
|
||||
except RuntimeError:
|
||||
return False
|
||||
return image_from_archive(str(path))['manifest_digest'] == digest
|
||||
return same_image(candidate, image_from_archive(str(path)))
|
||||
|
||||
if archive.exists():
|
||||
msg_info(translate('Verifying the image integrity...'))
|
||||
@@ -161,7 +161,7 @@ def update(vmid, acknowledge_external_data=False, proposal=None, keep_backup=Non
|
||||
changes = transaction.external_changes(record, config)
|
||||
transaction.preflight(record, desired, config)
|
||||
msg_ok(translate('Container checked'))
|
||||
current = record['observed']['image']['manifest_digest'] if operation == 'update' else None
|
||||
current = record['observed']['image'] if operation == 'update' else None
|
||||
archive, digest = resolve_archive(desired, config, current, lambda: transaction.require_backup_space(
|
||||
instances.location(instances.ROOT, vmid).parent, [vmid]))
|
||||
if archive is None:
|
||||
|
||||
@@ -17,10 +17,23 @@ find_snippet_storage() {
|
||||
fi
|
||||
storage=$(pvesm status --content snippets 2>/dev/null \
|
||||
| awk 'NR > 1 && $3 == "active" {print $1; exit}')
|
||||
[[ -n $storage ]] || die "No active storage supports snippets"
|
||||
if [[ -z $storage ]]; then
|
||||
# A new Proxmox installation accepts no snippets on any storage.
|
||||
enable_local_snippets || die "No active storage supports snippets"
|
||||
storage=local
|
||||
fi
|
||||
printf '%s' "$storage"
|
||||
}
|
||||
|
||||
enable_local_snippets() {
|
||||
local content
|
||||
content=$(pvesh get /storage/local --output-format json 2>/dev/null | jq -r '.content // empty')
|
||||
[[ -n $content ]] || return 1
|
||||
pvesm set local --content "${content},snippets" >/dev/null 2>&1 || return 1
|
||||
pvesm status --content snippets 2>/dev/null \
|
||||
| awk 'NR > 1 && $1 == "local" && $3 == "active" {found=1} END {exit !found}'
|
||||
}
|
||||
|
||||
install_hook() {
|
||||
local main_id=${1:?missing main VMID} source_config=${2:?missing lifecycle JSON}
|
||||
local storage hook_volume hook_path target_config
|
||||
|
||||
Reference in New Issue
Block a user