ProxMenux 1.2.6.2-beta: OCI containers in the Monitor, docs and fixes

OCI manager Apps
- App tab: containers installed from an OCI image are identified from their
  installation record; the application and image versions are shown and an
  update is detected by image digest; repository link; Refresh data.
- Updates tab for OCI containers: Update and Recreate run the same flow as the
  OCI menu in the Monitor terminal; the pre-update backup can be kept in a
  backup storage; scheduled image updates with an optional minimum age.
- Logs tab: console output of the application, kept on the host
  (lxc.console.logfile + logrotate) and followed live.
- The Proxmox console opens a shell (cmode: shell) when the image has one.
- A damaged image download is fetched again before failing.
- Multi-container applications open at their LAN address; volume mount
  points on block storage report their usage.

Monitor
- Proxmox notifications are delivered to a loopback-only HTTP listener when
  HTTPS is enabled, so they no longer fail certificate verification.
- Log persistence counts recurring patterns only; an ended burst is not
  reported as persistent and its warning clears on its own (#386).
- Proxmox notification config backups are deduplicated and capped at three.
- The update icon on the Apps page opens the container on its Updates tab.
- Version 1.2.6.2-beta and its release notes in every Monitor language.

Docs
- OCI manager Apps and Audit & Report rebuilt as per-page message files,
  with a new page for OCI containers in the Monitor.
- Seven pages fixed where rich-text tags were missing from t.rich.

Translations
- Spanish fixes across the OCI engine, the Monitor and the TUI menus.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
MacRimi
2026-09-25 21:51:12 +02:00
co-authored by Claude Opus 5.5
parent 386d33df6e
commit 4437a671d2
524 changed files with 14459 additions and 3841 deletions
+2
View File
@@ -136,6 +136,8 @@ cp "$SCRIPT_DIR/flask_proxmenux_routes.py" "$APP_DIR/usr/bin/" 2>/dev/null || ec
cp "$SCRIPT_DIR/post_install_versions.py" "$APP_DIR/usr/bin/" 2>/dev/null || echo "⚠️ post_install_versions.py not found"
cp "$SCRIPT_DIR/mount_monitor.py" "$APP_DIR/usr/bin/" 2>/dev/null || echo "⚠️ mount_monitor.py not found"
cp "$SCRIPT_DIR/lxc_mount_points.py" "$APP_DIR/usr/bin/" 2>/dev/null || echo "⚠️ lxc_mount_points.py not found"
cp "$SCRIPT_DIR/oci_console_logs.py" "$APP_DIR/usr/bin/" 2>/dev/null || echo "⚠️ oci_console_logs.py not found"
cp "$SCRIPT_DIR/oci_instance_info.py" "$APP_DIR/usr/bin/" 2>/dev/null || echo "⚠️ oci_instance_info.py not found"
cp "$SCRIPT_DIR/disk_temperature_history.py" "$APP_DIR/usr/bin/" 2>/dev/null || echo "⚠️ disk_temperature_history.py not found"
cp "$SCRIPT_DIR/smartctl_resolver.py" "$APP_DIR/usr/bin/" 2>/dev/null || echo "⚠️ smartctl_resolver.py not found"
cp "$SCRIPT_DIR/disk_identity.py" "$APP_DIR/usr/bin/" 2>/dev/null || echo "⚠️ disk_identity.py not found"
+52 -99
View File
@@ -865,108 +865,28 @@ _PVE_OUR_HEADERS = {
}
def _ssl_cert_hostname(cert_path: str) -> str:
"""Pull the most useful hostname out of an x509 cert.
Preference order: first DNS SAN → CN. Returns '' on any failure.
Used to build a webhook URL that won't fail PVE's TLS verification
(issue #239 — PVE has no `--insecure` flag and the user's ACME cert
is bound to a hostname, not to `127.0.0.1`).
"""
try:
import subprocess
out = subprocess.run(
['openssl', 'x509', '-in', cert_path, '-noout',
'-ext', 'subjectAltName', '-subject'],
capture_output=True, text=True, timeout=5,
)
if out.returncode != 0:
return ''
text = out.stdout or ''
# SAN line example: " DNS:pve.example.com, DNS:pve, IP Address:..."
import re
for line in text.splitlines():
m = re.search(r'DNS:([A-Za-z0-9.\-]+)', line)
if m:
return m.group(1)
# CN fallback. "subject= CN = pve.example.com" or "...CN=pve.example.com"
m = re.search(r'CN\s*=\s*([A-Za-z0-9.\-]+)', text)
if m:
return m.group(1)
except Exception:
pass
return ''
def _hostname_resolves_locally(hostname: str) -> bool:
"""True when `hostname` resolves to one of this host's own IPs.
Anti-misconfig: refuse to build a webhook URL that points
elsewhere (a stale DNS entry pointing at the previous host, a CN
that names a different node in the cluster, etc.). PVE delivers
webhooks from the same node, so the URL has to round-trip to
ourselves.
"""
try:
import socket
import ipaddress
target_ips = set()
for info in socket.getaddrinfo(hostname, None):
ip = info[4][0]
target_ips.add(ipaddress.ip_address(ip).compressed)
# Collect our own IPs from /proc/net/fib_trie isn't portable; use
# psutil if available, otherwise fall back to socket on the
# hostname itself.
local_ips = {'127.0.0.1', '::1'}
try:
import psutil
for _iface, addrs in psutil.net_if_addrs().items():
for a in addrs:
if a.family in (socket.AF_INET, socket.AF_INET6):
local_ips.add(ipaddress.ip_address(a.address.split('%')[0]).compressed)
except Exception:
pass
return bool(target_ips & local_ips)
except Exception:
return False
# With HTTPS enabled, PVE delivers its notifications to a plain-HTTP listener
# the Monitor opens on loopback only for this route (see flask_server). PVE's
# webhook client validates certificates against the system CA store, which
# does not hold a self-signed, ACME-staging or PVE-CA-signed Monitor
# certificate, so an https target fails with "unable to get local issuer
# certificate" on every notification.
WEBHOOK_LOOPBACK_PORT = 8009
def _pve_webhook_url() -> str:
"""Return the URL we register with PVE as our webhook target.
Three branches:
1. SSL off → http://127.0.0.1:8008 (always works, no cert).
2. SSL on + cert hostname extractable and resolves locally →
https://<cert-hostname>:8008. This is what PVE's TLS layer
actually validates against. Without this, PVE rejects the
self/ACME cert with "IP address mismatch" (issue #239).
3. SSL on but hostname extraction/check failed → fall back to
https://127.0.0.1:8008 and accept that the user may still
hit the cert-mismatch error. We log so the operator can
diagnose. Better than silently emitting a wrong URL.
"""
"""Return the URL we register with PVE as our webhook target: the
Monitor itself over HTTP when SSL is off, and its loopback-only HTTP
listener when SSL is on. Both stay on 127.0.0.1, so nothing crosses
the network and no certificate is involved."""
try:
from auth_manager import load_ssl_config
cfg = load_ssl_config() or {}
if not cfg.get('enabled'):
return 'http://127.0.0.1:8008/api/notifications/webhook'
cert_path = cfg.get('cert_path') or ''
if cert_path:
host = _ssl_cert_hostname(cert_path)
if host and _hostname_resolves_locally(host):
return f'https://{host}:8008/api/notifications/webhook'
if host:
print(
f"[ProxMenux] webhook URL fallback to 127.0.0.1: "
f"cert hostname '{host}' does not resolve to a local "
f"IP — PVE will likely report a TLS verification "
f"error. Fix by ensuring the FQDN resolves on this "
f"host (e.g. /etc/hosts entry)."
)
return 'https://127.0.0.1:8008/api/notifications/webhook'
if cfg.get('enabled'):
return f'http://127.0.0.1:{WEBHOOK_LOOPBACK_PORT}/api/notifications/webhook'
except Exception:
return 'http://127.0.0.1:8008/api/notifications/webhook'
pass
return 'http://127.0.0.1:8008/api/notifications/webhook'
# Backward-compat alias for callers that read this at import time. Most
@@ -988,15 +908,48 @@ def _pve_read_file(path):
return None, str(e)
_PVE_BACKUPS_KEPT = 3
def _pve_backup_file(path):
"""Create timestamped backup if file exists. Never fails fatally."""
"""Create timestamped backup if file exists. Never fails fatally.
The backups live next to the file, inside the cluster filesystem, whose
size is limited: a copy identical to the newest one is not written again
and only the newest few are kept. Restore uses the newest.
"""
import os, shutil
from datetime import datetime
try:
if os.path.exists(path):
if not os.path.exists(path):
return
directory, name = os.path.split(path)
prefix = f"{name}.proxmenux_backup_"
existing = sorted(
(os.path.join(directory, f) for f in os.listdir(directory) if f.startswith(prefix)),
key=os.path.getmtime, reverse=True,
)
with open(path, 'rb') as f:
current = f.read()
newest_matches = False
if existing:
try:
with open(existing[0], 'rb') as f:
newest_matches = f.read() == current
except OSError:
pass
if not newest_matches:
ts = datetime.now().strftime('%Y%m%d_%H%M%S')
backup = f"{path}.proxmenux_backup_{ts}"
shutil.copy2(path, backup)
shutil.copy2(path, os.path.join(directory, prefix + ts))
existing = sorted(
(os.path.join(directory, f) for f in os.listdir(directory) if f.startswith(prefix)),
key=os.path.getmtime, reverse=True,
)
for stale in existing[_PVE_BACKUPS_KEPT:]:
try:
os.remove(stale)
except OSError:
pass
except Exception:
pass
+143 -8
View File
@@ -678,6 +678,36 @@ def _warmup_lxc_ip_cache() -> int:
return count
def _lxc_isolated_ips(vmid):
"""Static addresses on interfaces that have no way out of the host.
lxc-info lists a container's addresses without saying which interface
carries each, so a container with a second network leg could be offered
at an address nobody can open: a ProxMenux stack wires its members through
a private bridge, and Paperless-ngx answered on 10.77.0.30 instead of its
LAN address. The interface a reader reaches is the one with a route out —
DHCP, or a static address with a gateway. A static address with no gateway
is local to the host, and is kept only as a last resort.
"""
try:
with open(f"/etc/pve/lxc/{int(vmid)}.conf", encoding="utf-8") as handle:
text = handle.read()
except (OSError, ValueError):
return set()
isolated = set()
for line in text.splitlines():
if line.startswith("["):
break # snapshot sections describe past states, not this one
match = re.match(r"net\d+:\s*(.*)", line)
if not match:
continue
options = dict(part.split("=", 1) for part in match.group(1).split(",") if "=" in part)
address = options.get("ip", "")
if address and address not in ("dhcp", "manual") and "gw" not in options:
isolated.add(address.split("/", 1)[0])
return isolated
def get_lxc_ip_from_lxc_info(vmid):
"""Get LXC IP addresses using lxc-info command (for DHCP containers)
Returns a dict with all IPs and classification"""
@@ -705,6 +735,9 @@ def get_lxc_ip_from_lxc_info(vmid):
else:
# Real network IPs (192.168.x.x, 10.x.x.x, etc.)
real_ips.append(ip)
isolated = _lxc_isolated_ips(vmid)
if isolated:
real_ips.sort(key=lambda ip: ip in isolated)
return {
'all_ips': ips,
@@ -3369,7 +3402,7 @@ def get_available_updates():
_system_info_cache['available_updates_stamp'] = stamp
return available_updates
# AGREGANDO FUNCIÓN PARA PARSEAR PROCESOS DE INTEL_GPU_TOP (SIN -J)
# Parse intel_gpu_top process output without -J.
def get_intel_gpu_processes_from_text():
"""Parse processes from intel_gpu_top text output (more reliable than JSON)"""
try:
@@ -3384,7 +3417,7 @@ def get_intel_gpu_processes_from_text():
bufsize=1
)
except FileNotFoundError:
# intel_gpu_top no está instalado, retornar lista vacía
# intel_gpu_top is not installed; return an empty list
return []
# Wait 2 seconds for intel_gpu_top to collect data
@@ -5519,7 +5552,7 @@ def _get_proxmox_storage_uncached():
for resource in resources:
node = resource.get('node', '')
# Filtrar solo storage del nodo local
# Keep only the local node storage
if node != local_node:
# print(f"[v0] Skipping storage {resource.get('storage')} from remote node: {node}")
pass
@@ -5538,7 +5571,7 @@ def _get_proxmox_storage_uncached():
pass
continue
# No filtrar storages no disponibles - mantenerlos para mostrar errores
# Keep unavailable storages so their errors stay visible
# Calcular porcentaje
percent = (used / total * 100) if total > 0 else 0.0
@@ -7832,7 +7865,7 @@ def get_detailed_gpu_info(gpu):
else:
# print(f"[v0] WARNING: No valid JSON objects found", flush=True)
pass
# CHANGE: Evitar bloqueo al leer stderr - usar communicate() con timeout
# communicate() with a timeout: reading stderr directly blocks
try:
# Use communicate() with timeout instead of read() to avoid blocking
_, stderr_output = process.communicate(timeout=0.5)
@@ -8312,8 +8345,7 @@ def get_detailed_gpu_info(gpu):
# print(f"[v0] Parsing fdinfo with {len(fdinfo)} entries", flush=True)
pass
# CHANGE: Corregir parseo de fdinfo con estructura anidada
# fdinfo es un diccionario donde las claves son los PIDs (como strings)
# fdinfo is nested: the keys are the PIDs, as strings
for pid_str, proc_data in fdinfo.items():
try:
process_info = {
@@ -13764,6 +13796,13 @@ def api_vm_apps_suggestions(vmid):
return jsonify({'error': str(e)}), 500
@app.route('/api/vms/<int:vmid>/apps/adguard-setup', methods=['GET'])
@require_auth
def api_vm_apps_adguard_setup(vmid):
import lxc_apps
return jsonify({'available': lxc_apps.oci_adguard_setup_available(vmid)})
@app.route('/api/vms/<int:vmid>/docker/inventory', methods=['GET'])
@require_auth
def api_vm_docker_inventory(vmid):
@@ -14113,6 +14152,8 @@ def _lxc_update_target_labels(
label = 'Applications'
elif target == 'docker-engine':
label = 'Docker Engine'
elif target == 'oci_image':
label = 'OCI image'
elif target.startswith('app:'):
app_item = apps.get(target.split(':', 1)[1]) or {}
label = str(app_item.get('name') or 'Application')
@@ -15077,7 +15118,7 @@ def api_vms_modal_cache_all():
return jsonify({'error': str(e)}), 500
# CHANGE: Modificar el endpoint para incluir la información completa de IPs
# The endpoint returns the complete IP information.
@app.route('/api/vms/<int:vmid>', methods=['GET'])
@require_auth
def get_vm_config(vmid):
@@ -15215,6 +15256,38 @@ def api_lxc_mount_points_runtime(vmid):
return jsonify({"ok": False, "error": str(e)}), 500
@app.route('/api/lxc/<int:vmid>/console-log', methods=['GET'])
@require_auth
def api_lxc_console_log(vmid):
"""Console output of a native OCI container. Never cached: the first
request returns the last `lines` lines, later ones pass back `offset`
and `inode` to receive only what was appended since."""
try:
import oci_console_logs
except ImportError as e:
return jsonify({"ok": False, "error": f"helper unavailable: {e}"}), 503
try:
lines = request.args.get('lines', default=200, type=int)
offset = request.args.get('offset', type=int)
inode = request.args.get('inode', type=int)
return jsonify(oci_console_logs.read(vmid, lines=lines, offset=offset, inode=inode))
except Exception as e:
return jsonify({"ok": False, "error": str(e)}), 500
@app.route('/api/lxc/<int:vmid>/oci-instance', methods=['GET'])
@require_auth
def api_lxc_oci_instance(vmid):
"""Whether the container was installed by OCI manager Apps, its stack,
host directories and pending operation, read from its record. Never
cached; it also says whether the console log exists (Logs tab)."""
try:
import oci_instance_info
return jsonify(oci_instance_info.info(vmid))
except Exception as e:
return jsonify({"ok": False, "error": str(e)}), 500
@app.route('/api/vms/<int:vmid>/logs', methods=['GET'])
@require_auth
def api_vm_logs(vmid):
@@ -22161,6 +22234,41 @@ def _run_scheduled_update(vmid: int, sched: dict) -> dict:
}
requested_target = sched.get("target") or "both"
if requested_targets == ['oci_image']:
# A container installed by OCI manager Apps is updated by replacing
# its image, through the same flow as the OCI menu, unattended.
command = ['bash', '/usr/local/share/proxmenux/scripts/oci/oci_manager_apps.sh',
'manage', str(vmid), '--action', 'update', '--unattended']
storage = (sched.get('backup_storage') or '').strip()
if sched.get('backup') and storage:
command += ['--keep-backup', storage]
if sched.get('acknowledge_external_data'):
command.append('--acknowledge-external-data')
delay = int(sched.get('release_delay_days') or 0)
if delay > 0:
command += ['--min-image-age-days', str(delay)]
try:
proc = subprocess.run(command, stdin=subprocess.DEVNULL, capture_output=True, text=True,
timeout=4 * 3600, env=dict(os.environ, TERM='dumb'))
except subprocess.TimeoutExpired:
reasons.append('the OCI image update did not finish in time')
return finish('failure', 'oci_image', [])
output = re.sub(r'\x1b\[[0-9;?]*[ -/]*[@-~]|\x1b[()][A-Z0-9]|\r', '', (proc.stdout or '') + (proc.stderr or ''))
_append_lxc_update_log(log_path, output)
if proc.returncode == 0:
return finish('success', 'oci_image', ['oci_image'])
if proc.returncode == 5:
deferred_targets.append('oci_image')
reasons.append(f'the new image is younger than {delay} days')
return finish('deferred', 'oci_image', [])
if proc.returncode == 4:
reasons.append('the container has changes made outside ProxMenux; review them in OCI manager Apps')
return finish('skipped', 'oci_image', [])
if proc.returncode == 3:
reasons.append('another OCI operation was running')
return finish('skipped', 'oci_image', [])
reasons.append('the OCI image update failed; the previous installation is restored when the update started')
return finish('failure', 'oci_image', [])
if not os.path.isfile(_APPLY_UPDATES_SCRIPT):
reasons.append('update runner is not installed')
return finish('skipped', requested_target, [])
@@ -22744,6 +22852,18 @@ if __name__ == '__main__':
print(f"[ProxMenux] SSL config error, falling back to HTTP: {e}")
ssl_ctx = None
# With HTTPS on, PVE cannot validate the Monitor certificate against the
# system CA store, so its webhook is delivered to this plain-HTTP
# listener instead. It binds 127.0.0.1 only and answers nothing but the
# webhook route.
from flask_notification_routes import WEBHOOK_LOOPBACK_PORT
def _webhook_loopback_app(environ, start_response):
if environ.get('PATH_INFO') == '/api/notifications/webhook':
return app(environ, start_response)
start_response('404 Not Found', [('Content-Type', 'text/plain')])
return [b'Not found']
# Use gevent for SSL+WebSocket support, or fallback to Flask dev server
gevent_available = False
ssl_loaded = False
@@ -22842,11 +22962,26 @@ if __name__ == '__main__':
ssl_context=ssl_context
)
gevent_available = True
try:
webhook_server = pywsgi.WSGIServer(
('127.0.0.1', WEBHOOK_LOOPBACK_PORT), _webhook_loopback_app, log=None)
webhook_server.start()
print(f"[ProxMenux] PVE webhook listener on 127.0.0.1:{WEBHOOK_LOOPBACK_PORT}")
except Exception as _e:
print(f"[ProxMenux] WARN: PVE webhook listener could not start on "
f"127.0.0.1:{WEBHOOK_LOOPBACK_PORT} ({_e}); PVE notifications will not arrive", flush=True)
server.serve_forever()
except ImportError as e:
print(f"[ProxMenux] gevent not available ({e})")
# Fallback: Flask dev server with SSL - flask-sock handles WebSockets
ssl_context = auth_manager.create_reloadable_ssl_context(ssl_cert, ssl_key)
try:
from werkzeug.serving import make_server
_webhook_srv = make_server('127.0.0.1', WEBHOOK_LOOPBACK_PORT, _webhook_loopback_app, threaded=True)
threading.Thread(target=_webhook_srv.serve_forever, daemon=True).start()
print(f"[ProxMenux] PVE webhook listener on 127.0.0.1:{WEBHOOK_LOOPBACK_PORT}")
except Exception as _e:
print(f"[ProxMenux] WARN: PVE webhook listener could not start ({_e})", flush=True)
print("[ProxMenux] Starting Flask server with SSL (using flask-sock for WebSockets)...")
app.run(host='::', port=8008, debug=False, ssl_context=ssl_context, threaded=True)
else:
+1 -1
View File
@@ -199,7 +199,7 @@ def search_command():
'command': stripped
})
# Resetear descripciones para el siguiente comando
# Reset the descriptions for the next command
current_description = []
return jsonify({
+32 -9
View File
@@ -4146,7 +4146,8 @@ class HealthMonitor:
New thresholds:
- CASCADE: ≥15 errors (increased from 10)
- SPIKE: ≥5 errors AND 4x increase (more restrictive)
- PERSISTENT: Same error in 3 consecutive checks
- PERSISTENT: the same pattern seen in at least 3 checks whose
occurrences span 15+ minutes, and still present in this check
"""
cache_key = 'logs_analysis'
current_time = time.time()
@@ -4193,6 +4194,8 @@ class HealthMonitor:
recent_patterns = defaultdict(int)
previous_patterns = defaultdict(int)
# Patterns present in this check; each counts one check once.
seen_this_check = set()
critical_errors_found = {} # To store unique critical error lines for persistence
for line in recent_lines:
@@ -4363,10 +4366,15 @@ class HealthMonitor:
else:
self.persistent_log_patterns[pattern] = {
'count': 1,
'checks': 0,
'first_seen': current_time,
'last_seen': current_time,
'sample': line.strip()[:200], # Original line for display
}
if pattern not in seen_this_check:
seen_this_check.add(pattern)
self.persistent_log_patterns[pattern]['checks'] = \
self.persistent_log_patterns[pattern].get('checks', 0) + 1
for line in previous_lines:
if not line.strip():
@@ -4409,10 +4417,23 @@ class HealthMonitor:
samples.append(clean[:120])
return samples
# A pattern that stopped appearing is forgotten before the
# evaluation, and its own warning (if it had one) is closed.
patterns_to_remove = [
p for p, data in self.persistent_log_patterns.items()
if current_time - data['last_seen'] > 1800
]
for pattern in patterns_to_remove:
del self.persistent_log_patterns[pattern]
# Persistent means recurring: present in this check, seen in at
# least three checks, and with occurrences spanning 15 minutes.
# A burst that ended is not persistent however long ago it began.
persistent_errors = {}
for pattern, data in self.persistent_log_patterns.items():
time_span = current_time - data['first_seen']
if data['count'] >= 3 and time_span >= 900: # 15 minutes
span = data['last_seen'] - data['first_seen']
if (pattern in seen_this_check and data.get('checks', 1) >= 3
and data['count'] >= 3 and span >= 900):
persistent_errors[pattern] = data['count']
# Record as warning if not already recorded
@@ -4437,12 +4458,14 @@ class HealthMonitor:
'dismissable': True, 'occurrences': data['count']}
)
patterns_to_remove = [
p for p, data in self.persistent_log_patterns.items()
if current_time - data['last_seen'] > 1800
]
for pattern in patterns_to_remove:
del self.persistent_log_patterns[pattern]
# Close the per-pattern warnings of patterns that are no
# longer persistent: the burst ended, or it was forgotten.
for pattern in list(self.persistent_log_patterns.keys()) + patterns_to_remove:
if pattern in persistent_errors:
continue
stale_key = f'log_persistent_{hashlib.md5(pattern.encode()).hexdigest()[:8]}'
if health_persistence.is_error_active(stale_key, category='logs'):
health_persistence.clear_error(stale_key)
# B5 fix: Cap size to prevent unbounded memory growth under high error load
MAX_LOG_PATTERNS = 500
+393 -7
View File
@@ -29,6 +29,7 @@ import datetime
import copy
import concurrent.futures
import hashlib
import ipaddress
import json
import os
import re
@@ -36,6 +37,7 @@ import signal
import shlex
import socket
import subprocess
import sys
import threading
import time
import urllib.error
@@ -62,7 +64,7 @@ _UPSTREAM_CACHE_TTL_SEC = 24 * 3600
_VALID_METHODS = ("dpkg", "apk", "file", "binary",
"python_dist", "docker_label", "docker_exec",
"command", "manual")
"command", "manual", "oci_image")
_DETECTOR_FIELDS = (
"package", "file_path", "file_regex", "binary_path", "binary_args",
"python_path", "distribution", "container_name", "label",
@@ -93,7 +95,7 @@ _VALID_UPSTREAM_TYPES = ("github", "http_json", "docker_hub")
# exposes and the freeform "custom" text field.
_VALID_SCHEDULE_TARGETS = ("os", "app", "both")
_SCHEDULE_TARGET_ID_RE = re.compile(
r"^(?:os|apps|app:[A-Za-z0-9_-]{1,64}|docker-engine|docker-(?:compose|container):[A-Za-z0-9][A-Za-z0-9_.-]{0,127}|docker-unit:[a-f0-9]{20})$"
r"^(?:os|apps|oci_image|app:[A-Za-z0-9_-]{1,64}|docker-engine|docker-(?:compose|container):[A-Za-z0-9][A-Za-z0-9_.-]{0,127}|docker-unit:[a-f0-9]{20})$"
)
_BULK_TARGET_ID_RE = re.compile(
r"^(?:os|app:[A-Za-z0-9_-]{1,64}|docker-engine|docker-unit:[a-f0-9]{20})$"
@@ -886,6 +888,10 @@ def validate_schedule(payload: Any) -> tuple[bool, Any]:
"backup_storage": backup_storage,
"restart": restart,
"release_delay_days": release_delay_days,
# Host directories of an OCI container are not reverted by its
# backup; a scheduled image update runs only when this was
# confirmed when the schedule was saved.
"acknowledge_external_data": bool(payload.get("acknowledge_external_data")),
}
# Preserve `last_run_at` / `last_run_status` when the caller sent
# them (typical when the scheduler writes back after firing);
@@ -951,7 +957,12 @@ def validate_config(payload: dict) -> tuple[bool, Any]:
if method:
conf["installed_via"] = method
if method in ("dpkg", "apk"):
if method == "oci_image":
# Everything it needs is in the installation record ProxMenux wrote:
# the image, the digest and the registry. There is no field to fill
# and no upstream to configure.
pass
elif method in ("dpkg", "apk"):
pkg = (payload.get("package") or "").strip()
if not pkg or not _PACKAGE_RE.match(pkg):
return _err("package is required (letters/digits/._+@:/ up to 127 chars)")
@@ -1345,6 +1356,15 @@ def detect_installed_version(vmid, config: dict) -> tuple[Optional[str], Optiona
# tag_regex is reused (backward-compat with older hints).
pattern = config.get("installed_regex") or config.get("tag_regex") or r"(\d+[.\d]+)"
if method == "oci_image":
result = _oci_image_versions(vmid, with_latest=False)
if result.get("error"):
return None, result["error"]
# An image that states no application version is still an image with
# a build date and a digest, which is what its updates are decided on.
return (result.get("installed_version")
or _oci_image_label(None, result.get("image_created"), result.get("installed_digest"))), None
if method == "dpkg":
rc, out, err = _pct_exec(vmid, ["dpkg-query", "-W", "-f=${Version}", config["package"]])
if rc != 0:
@@ -3767,9 +3787,14 @@ def partition_scheduled_release_targets(
# A delegated app has no updater of its own and never resolves a
# release date, so including it would hold the whole schedule back
# waiting for a date that will never arrive.
# An application ProxMenux installed from an OCI image is updated by
# replacing its image through the OCI engine's own transaction, not by
# a helper or a command inside the guest, so a scripted plan has
# nothing it could run for it.
if (not app_id or app.get("managed_oci_app_id")
or app.get("helper_slug") == "docker"
or app.get("update_via") == "docker"):
or app.get("update_via") == "docker"
or app.get("installed_via") == "oci_image"):
continue
if not select_all_apps and app_id not in selected_app_ids:
continue
@@ -3962,13 +3987,20 @@ def _app_update_notification_payload(vmid, app: dict) -> Optional[dict]:
return None
state = app.get("state") or {}
latest = state.get("latest_version")
installed = state.get("installed_version")
if app.get("installed_via") == "oci_image" and state.get("latest_digest"):
# What is published is an image. Naming it by version, build date and
# digest keeps two rebuilds of the same version from being taken for
# one notification, and tells the reader what actually changed.
latest = _oci_image_label(latest, state.get("latest_image_created"), state.get("latest_digest"))
installed = _oci_image_label(installed, state.get("image_created"), state.get("installed_digest"))
if not state.get("update_available") or not latest:
return None
return {
"vmid": int(vmid),
"ct_name": app.get("name") or f"CT-{vmid}",
"app_name": app.get("name") or "app",
"installed": state.get("installed_version") or "unknown",
"installed": installed or "unknown",
"latest": latest,
"app_id": str(app.get("id") or ""),
}
@@ -4291,6 +4323,28 @@ def check_app(
except (ValueError, TypeError):
pass
if app.get("installed_via") == "oci_image":
result = _oci_image_versions(vmid, known=state)
app["state"] = {
"installed_version": result.get("installed_version"),
"latest_version": result.get("latest_version"),
"latest_published_at": None,
"update_available": result.get("update_available"),
"error": result.get("error"),
"checked_at": _now_iso(),
"installed_digest": result.get("installed_digest"),
"latest_digest": result.get("latest_digest"),
"image_created": result.get("image_created"),
"latest_image_created": result.get("latest_image_created"),
"image_reference": result.get("image_reference"),
"image_repository": result.get("image_repository"),
}
sidecar["updated_at"] = _now_iso()
_write_sidecar(vmid, sidecar)
if notify and app["state"]["update_available"] and app["state"]["latest_version"]:
_fire_update_notification(vmid, app)
return sidecar
installed, inst_err, _healed = _detect_with_alt_healing(vmid, app)
# Trigger the upstream fetch when ANY upstream source is
# configured. The dispatcher inside `fetch_latest_upstream`
@@ -4511,6 +4565,12 @@ def _summarise_app(app: dict) -> dict:
# match this registered app against the CT's helper_slug and
# display its installed/upstream versions.
"helper_slug": app.get("helper_slug") or "",
# OCI image identity, for the Updates tab of an OCI container.
"image_reference": state.get("image_reference"),
"image_created": state.get("image_created"),
"installed_digest": state.get("installed_digest"),
"latest_image_created": state.get("latest_image_created"),
"latest_digest": state.get("latest_digest"),
}
@@ -5238,6 +5298,276 @@ def _helper_slug_meta(vmid) -> Optional[dict]:
return None
_OCI_INSTANCE_ROOT = "/usr/local/share/proxmenux/oci/instances"
_OCI_CATALOG_INDEX = "/usr/local/share/proxmenux/oci/engine/catalog/index.json"
_oci_catalog_cache: tuple[float, dict] | None = None
def _oci_catalog_icons() -> dict:
"""Current icon per catalog application, keyed by template id.
The installation record keeps a copy of the catalog entry as it stood on
the day of the install, which freezes the icon along with everything else.
Icons get corrected — most of the catalog used to point at a URL that
answered 404 — so the panel reads the catalog and keeps the record as the
fallback for an application the catalog no longer lists.
"""
global _oci_catalog_cache
now = time.time()
if _oci_catalog_cache and now - _oci_catalog_cache[0] < 600:
return _oci_catalog_cache[1]
icons: dict = {}
try:
with open(_OCI_CATALOG_INDEX, encoding="utf-8") as handle:
for item in (json.load(handle) or {}).get("applications", []):
icon = (item or {}).get("icon")
if not isinstance(icon, str) or not icon.startswith("http"):
continue
# An installation records the template id; the catalog is keyed
# by the application id and carries both.
for key in (item.get("template_id"), item.get("id")):
if key:
icons.setdefault(key, icon)
except (OSError, ValueError, TypeError):
pass
_oci_catalog_cache = (now, icons)
return icons
def _oci_localised(value) -> str:
"""A catalog_ui text field, which is either a string or a locale map."""
if isinstance(value, str):
return value.strip()
if isinstance(value, dict):
for key in ("en_US", "en", *sorted(value)):
text = value.get(key)
if isinstance(text, str) and text.strip():
return text.strip()
return ""
def _oci_name_from_image(reference: str | None) -> str:
"""A presentable name for an image that carries no catalog entry.
A stack member records the image it runs and the role it plays, not a
title: `nextcloud:latest` as the `application` of a Nextcloud stack.
"""
repository = str(reference or "").split("@", 1)[0]
basename = repository.rsplit("/", 1)[-1].rsplit(":", 1)[0].strip()
if not basename:
return ""
return " ".join(word.capitalize() for word in re.split(r"[-_.]+", basename) if word)
def _oci_instance_meta(vmid) -> Optional[dict]:
"""What ProxMenux itself recorded when it installed this container.
The sibling of `_helper_slug_meta`, and stronger evidence: a helper slug
is a hint read back out of the guest, while this is the contract the
installer wrote. It needs no `pct exec`, so it also answers for a stopped
container, and it cannot mistake an incidental binary for the application
— CT 152 runs Chromium and happens to have Docker inside, which the
runtime probe reported as the only candidate.
"""
try:
with open(f"{_OCI_INSTANCE_ROOT}/{int(vmid)}/oci-compose.json", encoding="utf-8") as handle:
record = json.load(handle)
except (OSError, ValueError, TypeError):
return None
if not isinstance(record, dict) or record.get("status") not in ("installed", "assembling"):
return None
template = record.get("template") or {}
contract = template.get("container_contract") or {}
image = contract.get("image") or {}
observed_image = (record.get("observed") or {}).get("image") or {}
# A stack member carries only its own image contract; the presentation
# belongs to the stack it is part of, which records its title, site,
# category and the endpoint the stack is reached on.
stack = record.get("stack") if isinstance(record.get("stack"), dict) else {}
stack_template = stack.get("template") if isinstance(stack.get("template"), dict) else {}
role = str((record.get("stack_member") or {}).get("name") or "").strip()
member = record.get("stack_member") or {}
is_primary = not stack_template or member.get("primary_vmid") in (None, record.get("vmid"))
ui = template.get("catalog_ui") or stack_template.get("catalog_ui") or {}
if not template.get("catalog_ui") and stack_template and is_primary:
# Only the member the stack is reached on inherits the stack endpoint.
# The cache and the database of a Nextcloud stack do not answer on its
# web port, and offering it would register a service that is not there.
template = {**stack_template, "id": template.get("id") or stack_template.get("id")}
elif not template.get("catalog_ui") and stack_template:
ui = {key: value for key, value in ui.items()
if key in ("category", "category_label")}
# `container_contract.ports` lists every port the image exposes; the
# endpoint is the one the application is actually reached on, with its
# scheme. Chromium exposes 3000 and 3001 and serves on 3001 over https.
endpoint = next((e for e in (template.get("first_run") or {}).get("endpoints") or []
if isinstance(e, dict) and e.get("port")), None) or ui.get("launch") or {}
ports = []
for entry in contract.get("ports") or []:
try:
port = int((entry or {}).get("container_port"))
except (TypeError, ValueError):
continue
if 1 <= port <= 65535 and port not in ports:
ports.append(port)
template_id = str(template.get("id") or "").strip()
catalog_icons = _oci_catalog_icons()
logo = catalog_icons.get(template_id) or ""
if not logo:
# A stack or a one-off image has no catalog entry of its own, but the
# image it runs usually does: the Nextcloud stack wears Nextcloud's.
repository = str(image.get("reference") or "").split("@", 1)[0]
basename = repository.rsplit("/", 1)[-1].rsplit(":", 1)[0].strip().lower()
logo = catalog_icons.get(basename) or ""
if not logo:
logo = ui.get("icon") if isinstance(ui.get("icon"), str) else ""
website = ui.get("website") if isinstance(ui.get("website"), str) else ""
return {
"template_id": template_id or None,
"name": (_oci_localised(ui.get("title"))
or _oci_name_from_image(image.get("reference"))
or str(template.get("id") or "").strip() or None),
"role": role or None,
"logo": logo if logo.startswith(("http://", "https://")) else "",
"website": website if website.startswith(("http://", "https://")) else "",
"category": str(ui.get("category") or "").strip() or None,
"category_label": str(ui.get("category_label") or "").strip() or None,
"image_reference": str(image.get("reference") or "").strip() or None,
# The repository of the image — its GitHub project, or its page on the
# registry for an official image — which is where its updates come from.
"repository": (str(ui.get("repository") or "").strip()
or str((template.get("source") or {}).get("repository") or "").strip() or None),
"endpoint_port": endpoint.get("port") if isinstance(endpoint.get("port"), int) else None,
"endpoint_scheme": str(endpoint.get("scheme") or "").strip().lower() or None,
"endpoint_path": str(endpoint.get("path") or "").strip() or None,
"ports": ports,
# The exact image this container was created from. Its digest is what
# an update is decided on; the version label is only for reading.
"installed_digest": str(observed_image.get("manifest_digest") or "").strip() or None,
"architecture": str(observed_image.get("architecture") or "").strip() or None,
}
def oci_adguard_setup_available(vmid) -> bool:
"""Probe only this OCI application's setup endpoint, without caching it."""
meta = _oci_instance_meta(vmid)
if not meta or meta.get('template_id') != 'image-adguard-home':
return False
try:
result = subprocess.run(['lxc-info', '-n', str(int(vmid)), '-iH'],
capture_output=True, text=True, timeout=3, check=True)
ip = next(str(ipaddress.IPv4Address(value.strip()))
for value in result.stdout.splitlines()
if value.strip() and ipaddress.ip_address(value.strip()).version == 4)
with socket.create_connection((ip, 3000), timeout=1) as connection:
connection.settimeout(1)
connection.sendall(b'GET / HTTP/1.0\r\nHost: localhost\r\n\r\n')
return connection.recv(32).startswith(b'HTTP/')
except (OSError, ValueError, StopIteration, subprocess.SubprocessError):
return False
_OCI_REMOTE_DIR = "/usr/local/share/proxmenux/oci/engine/remote"
def _oci_state_module():
"""The OCI engine's own reader of image versions.
Reused rather than copied: it resolves the platform manifest, reads the
version the image states in its environment before the label it may
have inherited from its base, and it is the same code the engine
installs and updates with, so the panel and the updater cannot disagree
about what an image is.
"""
if _OCI_REMOTE_DIR not in sys.path:
sys.path.insert(0, _OCI_REMOTE_DIR)
import oci_installation_state
return oci_installation_state
def _oci_repository(reference: str) -> str:
repository = reference.split("@", 1)[0]
if ":" in repository.rsplit("/", 1)[-1]:
repository = repository.rsplit(":", 1)[0]
return repository
def _oci_resolve(module, reference: str, architecture: str) -> dict:
# One retry: an anonymous registry answers the occasional request with an
# error that the next one does not repeat.
try:
return module.resolve_candidate(reference, architecture)
except RuntimeError:
time.sleep(2)
return module.resolve_candidate(reference, architecture)
def _oci_image_versions(vmid, known: Optional[dict] = None, with_latest: bool = True) -> dict:
"""Installed and published version of a container ProxMenux installed.
Nothing is inferred. The installation record names the image and the
exact digest the container was created from; the registry says which
digest the same tag points at today. An update is available when those
two differ — the version strings only say which one it is, and they are
not compared, because a rebuild can keep its number and a build id such
as ``b1ee1dc8-ls55`` has no order to compare.
The installed version is read by digest, which never changes, so a
previous answer for the same digest is reused instead of asked again.
"""
meta = _oci_instance_meta(vmid)
if not meta or not meta.get("image_reference") or not meta.get("installed_digest"):
return {"error": "no OCI installation record for this container"}
reference = meta["image_reference"]
architecture = meta.get("architecture") or "amd64"
installed_digest = meta["installed_digest"]
result: dict = {"installed_digest": installed_digest, "image_reference": reference,
"image_repository": meta.get("repository")}
try:
module = _oci_state_module()
except Exception as exc:
return {**result, "error": f"OCI engine unavailable: {exc}"}
known = known or {}
if (known.get("installed_digest") == installed_digest
and (known.get("installed_version") or known.get("image_created"))):
result["installed_version"] = known.get("installed_version")
result["image_created"] = known.get("image_created")
else:
try:
installed = _oci_resolve(
module, f"{_oci_repository(reference)}@{installed_digest}", architecture)
result["installed_version"] = installed.get("version")
result["image_created"] = installed.get("created")
except Exception as exc:
return {**result, "error": f"could not read the installed image: {exc}"}
if not with_latest:
return result
try:
latest = _oci_resolve(module, reference, architecture)
except Exception as exc:
return {**result, "error": f"could not read {reference} from its registry: {exc}"}
latest_digest = latest.get("manifest_digest")
# The image decides. An application whose version did not move can still
# have a new image — a rebuild on a patched base — and that is an update
# for a container whose application only changes when its image does.
replaced = bool(latest_digest) and latest_digest != installed_digest
result.update(latest_digest=latest_digest, latest_version=latest.get("version"),
latest_image_created=latest.get("created"), update_available=replaced)
return result
def _oci_image_label(version: Optional[str], created: Optional[str], digest: Optional[str]) -> str:
"""One image, as a line a reader can compare: version, build date, digest."""
parts = [version] if version else []
if created:
parts.append(str(created)[:10])
if digest:
parts.append(str(digest).split(":", 1)[-1][:8])
return " · ".join(parts) or "unknown"
def _catalog_lookup(slug: str) -> Optional[dict]:
"""Fetch the community-scripts catalog entry for a slug.
Returns {name, updateable, default_port, logo} or None.
@@ -5292,6 +5622,7 @@ def get_suggestions(vmid, force: bool = False) -> dict:
if p in _KNOWN_WEB_PORTS:
web_hint = "/"
break
oci_meta = _oci_instance_meta(vmid)
meta = _helper_slug_meta(vmid) or {}
slug = meta.get("slug")
# Suppress base-OS helper slugs from the suggestion pipeline.
@@ -5318,7 +5649,10 @@ def get_suggestions(vmid, force: bool = False) -> dict:
detector.get("installed_via") in ("docker_label", "docker_exec")
for detector in primary_matches
)
if "docker" in detected_map and primary_is_docker_workload:
# An OCI container installed by ProxMenux is the application named in its
# own record. Promoting a probed binary over that would offer Docker for a
# Chromium container just because the image ships a docker client.
if "docker" in detected_map and primary_is_docker_workload and not oci_meta:
slug = "docker"
meta = {"slug": "docker", "name": "Docker"}
# Tracking hint pipeline: catalog + curated hints merged.
@@ -5560,6 +5894,47 @@ def get_suggestions(vmid, force: bool = False) -> dict:
"tracking_suggestion": det_tracking,
})
if oci_meta:
# The record states what this container runs, so a probe finding is
# noise: CT 152 runs Chromium and ships a docker client, and offering
# to register Docker there invites the user to track the wrong thing.
extras = []
# Recorded facts replace every probed guess: the name, the logo, the
# site and the endpoint the application is served on. The version is
# not filled here — the image has no detector among the guest-side
# methods — so the entry registers without version tracking until the
# OCI detector exists.
name_sug = oci_meta["name"] or name_sug
logo_url = oci_meta["logo"] or logo_url
category_suggestion = oci_meta["category_label"] or suggest_category_for(slug)
if oci_meta["endpoint_port"]:
default_ports = [oci_meta["endpoint_port"]]
ports = [oci_meta["endpoint_port"]] + [p for p in (oci_meta["ports"] or ports)
if p != oci_meta["endpoint_port"]]
elif oci_meta["ports"]:
default_ports = list(oci_meta["ports"])
web_hint = oci_meta["endpoint_path"] or web_hint
if oci_meta["template_id"] == "image-adguard-home":
# The first-run endpoint disappears once setup switches to :80.
default_ports = [80]
ports = [80, 3000]
# Version tracking comes with the registration. Its updates are
# decided by the image, which always has a build date and a digest,
# so it applies even to an image that states no application version.
versions = _oci_image_versions(vmid, with_latest=False)
if not versions.get("error"):
tracking = {
"installed_via": "oci_image",
"detected_version": versions.get("installed_version") or _oci_image_label(
None, versions.get("image_created"), versions.get("installed_digest")),
"detector_verified": True,
"detector_source": "oci_instance_record",
"logo": logo_url or None,
"website": oci_meta["website"] or None,
}
else:
category_suggestion = suggest_category_for(slug)
return {
"name_suggestion": name_sug,
"helper_slug": slug,
@@ -5568,9 +5943,20 @@ def get_suggestions(vmid, force: bool = False) -> dict:
"tracking_suggestion": tracking,
"default_ports": default_ports,
"logo_url": logo_url or None,
# Identity of a ProxMenux OCI install: the scheme the endpoint is
# served on, the image it was created from, and the upstream source.
"oci_instance": {
"template_id": oci_meta["template_id"],
"image_reference": oci_meta["image_reference"],
"repository": oci_meta["repository"],
"scheme": oci_meta["endpoint_scheme"],
"port": oci_meta["endpoint_port"],
"path": oci_meta["endpoint_path"],
"website": oci_meta["website"],
} if oci_meta else None,
# Categoría preset for the primary detection — same lookup as
# get_catalog_entry so the Register button pre-selects it.
"category": suggest_category_for(slug),
"category": category_suggestion,
"extras": extras,
"docker_workloads": sorted(docker_workloads, key=lambda item: item["name"].lower()),
"docker_web_links": docker_web_links,
+46 -11
View File
@@ -304,21 +304,56 @@ def _df_via_host_pid(host_pid: str, ct_target: str) -> dict[str, Optional[int]]:
["df", "-B1", "--output=size,used,avail", full],
capture_output=True, text=True, timeout=_STAT_TIMEOUT,
)
except subprocess.TimeoutExpired:
# A filesystem that does not answer df will not answer statfs
# either; asking again would only double the wait.
return empty
except OSError:
return empty
if proc.returncode == 0:
lines = [ln for ln in proc.stdout.strip().splitlines() if ln.strip()]
parts = lines[-1].split() if len(lines) >= 2 else []
if len(parts) >= 3:
try:
return {
"total_bytes": int(parts[0]),
"used_bytes": int(parts[1]),
"available_bytes": int(parts[2]),
}
except ValueError:
pass
return _statfs_via_host(full)
def _statfs_via_host(full: str) -> dict[str, Optional[int]]:
"""Capacity of a path read from its filesystem rather than the mount table.
df names a path by finding the mount that holds it in the host's own
table, and a volume of a native OCI container is not in that table: df
answers "no file systems processed" and the tab showed no usage for any
of its volumes. statfs asks the filesystem of the path directly and
returns the same three figures df prints. It is only reached when df
completes without an answer, so every mount df can measure keeps the
value it had.
"""
empty = {"total_bytes": None, "used_bytes": None, "available_bytes": None}
try:
proc = subprocess.run(
["stat", "-f", "-c", "%S %b %f %a", full],
capture_output=True, text=True, timeout=_STAT_TIMEOUT,
)
if proc.returncode != 0:
return empty
lines = [ln for ln in proc.stdout.strip().splitlines() if ln.strip()]
if len(lines) < 2:
return empty
parts = lines[-1].split()
if len(parts) < 3:
return empty
return {
"total_bytes": int(parts[0]),
"used_bytes": int(parts[1]),
"available_bytes": int(parts[2]),
}
block, total, free, available = (int(value) for value in proc.stdout.split())
except (subprocess.TimeoutExpired, OSError, ValueError):
return empty
if block <= 0 or total <= 0:
return empty
return {
"total_bytes": total * block,
"used_bytes": (total - free) * block,
"available_bytes": available * block,
}
def _df_via_pct_exec(vmid: str, ct_target: str,
+1 -2
View File
@@ -97,8 +97,7 @@ class RateLimiter:
# Counter of events dropped while over the rate limit. Surfaced via
# `consume_drop_count()` so the dispatch loop can periodically log
# "X events suppressed by rate-limit" instead of letting them
# disappear silently. Audit Tier 6 — `RateLimiter` descarta
# silenciosamente eventos sobre el límite.
# disappear silently.
self._dropped: int = 0
def allow(self) -> bool:
+1 -2
View File
@@ -2502,8 +2502,7 @@ class TaskWatcher:
# Manual starts (onboot=0) within the grace period also bypass the
# aggregator: a user manually starting a VM right after boot wants
# the individual confirmation, not their action silently rolled into
# the autostart summary. Audit Tier 6 — `system_startup` aggregation
# puede tragar VM starts manuales del usuario durante grace period.
# the autostart summary.
_STARTUP_EVENTS = {'vm_start', 'ct_start'}
if event_type in _STARTUP_EVENTS and not is_error and not is_warning:
if _shared_state.is_startup_period():
+1 -2
View File
@@ -1353,8 +1353,7 @@ TEMPLATES = {
# `system-mail` event, and the Monitor forwards it to every enabled
# channel. Most operators want smartd alerts but NOT noisy cron
# output — without a visible toggle the only fix is editing
# /etc/aliases or removing MAILTO from the cron job. Audit Tier 6
# — `system_mail` toggle no visible en UI / reportado por usuario.
# /etc/aliases or removing MAILTO from the cron job.
},
'apt_listchanges': {
'title': '{hostname}: {pve_title}',
+132
View File
@@ -0,0 +1,132 @@
"""Console output of a native OCI container, read from the host.
ProxMenux creates every OCI container with `lxc.console.logfile` pointing at
`/var/log/proxmenux/oci/<vmid>.console.log`: liblxc copies the stdout and
stderr of the image's entrypoint there, the way `docker logs` keeps them. The
file already exists on the host, so nothing here runs inside the container.
The path is derived from the VMID and confirmed against the container's own
configuration; no caller can name a file. Reads are bounded — the last lines
on first open, then only what was appended since the offset the viewer holds —
and a file that got shorter or was replaced (logrotate's copytruncate, a
container rebuilt by an update) resets the viewer instead of returning garbage.
"""
from __future__ import annotations
import os
import re
LOG_DIR = "/var/log/proxmenux/oci"
TAIL_READ_LIMIT = 2 * 1024 * 1024 # bytes scanned to find the last lines
FOLLOW_READ_LIMIT = 512 * 1024 # bytes returned per follow request
MAX_LINES = 2000
# CSI sequences (colours, cursor), OSC sequences (window titles) and the lone
# ESC-letter codes some programs emit. Stripped, never rendered: the text goes
# to the page as text.
_ANSI_RE = re.compile(r"\x1b\[[0-9;?]*[ -/]*[@-~]|\x1b\][^\x07\x1b]*(?:\x07|\x1b\\)|\x1b[@-Z\\-_]")
_CONTROL_RE = re.compile(r"[\x00-\x08\x0b\x0c\x0e-\x1f\x7f]")
def log_path(vmid: int) -> str:
return os.path.join(LOG_DIR, f"{int(vmid)}.console.log")
def configured_log(vmid: int) -> str | None:
"""The console log the container declares, when it is the ProxMenux one."""
try:
with open(f"/etc/pve/lxc/{int(vmid)}.conf", encoding="utf-8") as handle:
for line in handle:
if line.startswith("["):
break
if line.startswith("lxc.console.logfile:"):
value = line.split(":", 1)[1].strip()
return value if value == log_path(vmid) else None
except (OSError, ValueError):
return None
return None
def terminal_state(vmid: int) -> dict:
"""Whether the Proxmox console of the container opens a shell."""
mode = "tty"
try:
with open(f"/etc/pve/lxc/{int(vmid)}.conf", encoding="utf-8") as handle:
for line in handle:
if line.startswith("["):
break
if line.startswith("cmode:"):
mode = line.split(":", 1)[1].strip()
except (OSError, ValueError):
pass
return {"cmode": mode, "shell": mode == "shell"}
def clean(text: str) -> list[str]:
"""Readable lines from raw console output.
CRLF becomes a line end. A lone carriage return is how a progress bar
redraws itself in place, so only what was written after the last one on a
line is kept — the state the terminal would have shown.
"""
text = _ANSI_RE.sub("", text).replace("\r\n", "\n")
lines = []
for raw in text.split("\n"):
if "\r" in raw:
raw = raw.rsplit("\r", 1)[-1]
lines.append(_CONTROL_RE.sub("", raw))
return lines
def _read_tail(handle, size: int, lines: int) -> tuple[list[str], bool]:
start = max(0, size - TAIL_READ_LIMIT)
handle.seek(start)
data = handle.read(size - start).decode("utf-8", errors="replace")
result = clean(data)
if start > 0 and result:
result = result[1:] # the first line was cut by the read window
if result and result[-1] == "":
result = result[:-1]
head_cut = start > 0 or len(result) > lines
return result[-lines:], head_cut
def read(vmid: int, lines: int = 200, offset: int | None = None, inode: int | None = None) -> dict:
"""The last `lines` lines, or what was appended after `offset`.
`lines=0` reads nothing: it only says whether the container has a
console log, which is what decides if the viewer is offered at all.
"""
lines = max(0, min(int(lines), MAX_LINES))
path = configured_log(vmid)
base = {"ok": True, "vmid": int(vmid), **terminal_state(vmid)}
if not path:
return {**base, "enabled": False, "lines": [], "size": 0, "offset": 0, "inode": None}
try:
stat = os.stat(path)
except OSError:
return {**base, "enabled": True, "path": path, "lines": [], "size": 0, "offset": 0, "inode": None}
size = stat.st_size
base.update(enabled=True, path=path, size=size, inode=stat.st_ino)
if lines == 0 and offset is None:
return {**base, "lines": [], "offset": size}
with open(path, "rb") as handle:
replaced = inode is not None and inode != stat.st_ino
if offset is None or replaced or offset > size:
# First open, or the file the viewer followed is gone: rotated by
# copytruncate, or recreated with the container.
result, head_cut = _read_tail(handle, size, max(lines, 1))
return {**base, "lines": result, "offset": size, "reset": offset is not None,
"head_truncated": head_cut}
handle.seek(offset)
chunk = handle.read(min(size - offset, FOLLOW_READ_LIMIT))
# Only whole lines are returned; an unfinished last line is left for the
# next read, where it arrives complete.
cut = chunk.rfind(b"\n")
if cut < 0:
return {**base, "lines": [], "offset": offset, "reset": False}
complete = chunk[:cut + 1]
result = clean(complete.decode("utf-8", errors="replace"))
if result and result[-1] == "":
result = result[:-1]
return {**base, "lines": result, "offset": offset + len(complete), "reset": False}
+62
View File
@@ -0,0 +1,62 @@
"""What the VM & LXC modal needs to know about an OCI instance.
Read-only view of the installation record OCI manager Apps keeps for every
container it created: whether the container is one, whether it belongs to a
multi-container application, whether it uses host directories (which its
backup does not revert) and whether an operation is pending. Nothing here
changes the record or runs inside the container.
"""
from __future__ import annotations
import json
import os
import oci_console_logs
ROOT = "/usr/local/share/proxmenux/oci/instances"
def _record(vmid: int) -> dict | None:
path = os.path.join(ROOT, str(int(vmid)), "oci-compose.json")
try:
with open(path, encoding="utf-8") as handle:
record = json.load(handle)
except (OSError, ValueError):
return None
return record if isinstance(record, dict) else None
def _host_dirs(record: dict) -> bool:
mounts = (record.get("deployment") or {}).get("mounts") or []
return any(isinstance(m, dict) and m.get("type") == "host-bind" for m in mounts)
def info(vmid: int) -> dict:
vmid = int(vmid)
result = {
"vmid": vmid,
"oci_instance": False,
"console_log": oci_console_logs.configured_log(vmid) is not None,
"stack": False,
"primary_vmid": vmid,
"members": [],
"host_directories": False,
"pending": False,
}
record = _record(vmid)
if record is None:
return result
primary_id = int((record.get("stack_member") or {}).get("primary_vmid") or vmid)
primary = record if primary_id == vmid else (_record(primary_id) or {})
members = [int(m["vmid"]) for m in (primary.get("stack") or {}).get("members") or []
if isinstance(m, dict) and str(m.get("vmid", "")).isdigit()]
records = [record] if not members else [r for r in (_record(m) for m in members) if r]
result.update(
oci_instance=True,
stack=bool(members) or bool(record.get("stack_member")),
primary_vmid=primary_id,
members=members,
host_directories=any(_host_dirs(r) for r in records),
pending=bool(record.get("pending_transaction") or primary.get("pending_stack_transaction")),
)
return result
@@ -0,0 +1,217 @@
import importlib.util
import json
import shutil
import socket
import ssl
import subprocess
import tempfile
import threading
import unittest
from pathlib import Path
from unittest import mock
MODULE_PATH = Path(__file__).resolve().parents[1] / "auth_manager.py"
SPEC = importlib.util.spec_from_file_location("auth_manager_ssl_under_test", MODULE_PATH)
auth_manager = importlib.util.module_from_spec(SPEC)
SPEC.loader.exec_module(auth_manager)
class _FakeSslSocket:
def __init__(self, context):
self.context = context
class ProxmoxCertificateHotReloadTests(unittest.TestCase):
@classmethod
def setUpClass(cls):
if shutil.which("openssl") is None:
raise unittest.SkipTest("openssl is required for TLS fixture generation")
cls.fixture_dir = tempfile.TemporaryDirectory()
fixture_path = Path(cls.fixture_dir.name)
cls.pairs = []
for name in ("original", "renewed"):
cert_path = fixture_path / f"{name}.pem"
key_path = fixture_path / f"{name}.key"
subprocess.run(
[
"openssl", "req", "-x509", "-newkey", "rsa:2048",
"-nodes", "-days", "1", "-subj", f"/CN={name}.test",
"-keyout", str(key_path), "-out", str(cert_path),
],
check=True,
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
)
cls.pairs.append((cert_path, key_path))
@classmethod
def tearDownClass(cls):
cls.fixture_dir.cleanup()
def setUp(self):
self.temp_dir = tempfile.TemporaryDirectory()
self.addCleanup(self.temp_dir.cleanup)
temp_path = Path(self.temp_dir.name)
self.active_cert = temp_path / "pveproxy-ssl.pem"
self.active_key = temp_path / "pveproxy-ssl.key"
self.ssl_config = temp_path / "ssl_config.json"
self._install_pair(0)
self._write_config("proxmox")
self.patch = mock.patch.multiple(
auth_manager,
SSL_CONFIG_FILE=self.ssl_config,
PROXMOX_CUSTOM_CERT_PATH=str(self.active_cert),
PROXMOX_CUSTOM_KEY_PATH=str(self.active_key),
PROXMOX_CERT_PATH=str(temp_path / "missing-pve-ssl.pem"),
PROXMOX_KEY_PATH=str(temp_path / "missing-pve-ssl.key"),
)
self.patch.start()
self.addCleanup(self.patch.stop)
self.addCleanup(self._reset_runtime)
def _reset_runtime(self):
with auth_manager._SSL_RUNTIME_LOCK:
auth_manager._SSL_RUNTIME_CONTEXT = None
auth_manager._SSL_RUNTIME_FINGERPRINT = ""
auth_manager._SSL_RUNTIME_CERT_PATH = ""
auth_manager._SSL_RUNTIME_KEY_PATH = ""
auth_manager._SSL_RUNTIME_SOURCE = "none"
auth_manager._SSL_RUNTIME_LAST_REFRESH_ERROR = ""
def _install_pair(self, index):
cert_path, key_path = self.pairs[index]
shutil.copyfile(cert_path, self.active_cert)
shutil.copyfile(key_path, self.active_key)
def _write_config(self, source):
self.ssl_config.write_text(json.dumps({
"enabled": True,
"cert_path": str(self.active_cert),
"key_path": str(self.active_key),
"source": source,
}))
def _create_context(self):
return auth_manager.create_reloadable_ssl_context(
str(self.active_cert), str(self.active_key)
)
def test_unchanged_pair_does_not_rebuild_the_context(self):
context = self._create_context()
ssl_socket = _FakeSslSocket(context)
with mock.patch.object(
auth_manager,
"reload_server_ssl_context",
wraps=auth_manager.reload_server_ssl_context,
) as reload_mock:
context.sni_callback(ssl_socket, "proxmenux.test", context)
reload_mock.assert_not_called()
self.assertIs(ssl_socket.context, context)
def test_valid_renewed_pair_is_activated_during_the_handshake(self):
context = self._create_context()
previous_fingerprint = auth_manager._SSL_RUNTIME_FINGERPRINT
self._install_pair(1)
ssl_socket = _FakeSslSocket(context)
context.sni_callback(ssl_socket, "proxmenux.test", context)
self.assertNotEqual(previous_fingerprint, auth_manager._SSL_RUNTIME_FINGERPRINT)
self.assertIs(ssl_socket.context, auth_manager._SSL_RUNTIME_CONTEXT)
self.assertIsNot(ssl_socket.context, context)
self.assertEqual(auth_manager._SSL_RUNTIME_LAST_REFRESH_ERROR, "")
def test_mismatched_pair_keeps_the_previous_context(self):
context = self._create_context()
previous_fingerprint = auth_manager._SSL_RUNTIME_FINGERPRINT
shutil.copyfile(self.pairs[1][0], self.active_cert)
ssl_socket = _FakeSslSocket(context)
context.sni_callback(ssl_socket, "proxmenux.test", context)
self.assertEqual(previous_fingerprint, auth_manager._SSL_RUNTIME_FINGERPRINT)
self.assertIs(auth_manager._SSL_RUNTIME_CONTEXT, context)
self.assertIs(ssl_socket.context, context)
self.assertTrue(auth_manager._SSL_RUNTIME_LAST_REFRESH_ERROR)
def test_custom_certificate_source_is_not_examined_automatically(self):
self._write_config("custom")
context = self._create_context()
previous_fingerprint = auth_manager._SSL_RUNTIME_FINGERPRINT
self._install_pair(1)
ssl_socket = _FakeSslSocket(context)
context.sni_callback(ssl_socket, "proxmenux.test", context)
self.assertEqual(previous_fingerprint, auth_manager._SSL_RUNTIME_FINGERPRINT)
self.assertIs(ssl_socket.context, context)
def test_first_real_tls_connection_receives_the_renewed_certificate(self):
server_context = self._create_context()
self._install_pair(1)
server_socket, client_socket = socket.socketpair()
server_socket.settimeout(5)
client_socket.settimeout(5)
server_error = []
def serve_once():
try:
with server_context.wrap_socket(server_socket, server_side=True) as tls_socket:
tls_socket.recv(1)
except Exception as error: # pragma: no cover - asserted below
server_error.append(error)
thread = threading.Thread(target=serve_once)
thread.start()
client_context = ssl.SSLContext(ssl.PROTOCOL_TLS_CLIENT)
client_context.check_hostname = False
client_context.verify_mode = ssl.CERT_NONE
with client_context.wrap_socket(
client_socket, server_hostname="proxmenux.test"
) as tls_client:
received_der = tls_client.getpeercert(binary_form=True)
tls_client.sendall(b"x")
thread.join(timeout=5)
self.assertFalse(thread.is_alive())
self.assertEqual(server_error, [])
renewed_pem = self.pairs[1][0].read_text()
self.assertEqual(received_der, ssl.PEM_cert_to_DER_cert(renewed_pem))
def test_tls_connection_without_sni_also_receives_the_renewed_certificate(self):
server_context = self._create_context()
self._install_pair(1)
server_socket, client_socket = socket.socketpair()
server_socket.settimeout(5)
client_socket.settimeout(5)
server_error = []
def serve_once():
try:
with server_context.wrap_socket(server_socket, server_side=True) as tls_socket:
tls_socket.recv(1)
except Exception as error: # pragma: no cover - asserted below
server_error.append(error)
thread = threading.Thread(target=serve_once)
thread.start()
client_context = ssl.SSLContext(ssl.PROTOCOL_TLS_CLIENT)
client_context.check_hostname = False
client_context.verify_mode = ssl.CERT_NONE
with client_context.wrap_socket(client_socket) as tls_client:
received_der = tls_client.getpeercert(binary_form=True)
tls_client.sendall(b"x")
thread.join(timeout=5)
self.assertFalse(thread.is_alive())
self.assertEqual(server_error, [])
renewed_pem = self.pairs[1][0].read_text()
self.assertEqual(received_der, ssl.PEM_cert_to_DER_cert(renewed_pem))
if __name__ == "__main__":
unittest.main()
@@ -0,0 +1,69 @@
import sys
import tempfile
import unittest
from pathlib import Path
from unittest import mock
SCRIPTS_DIR = Path(__file__).resolve().parents[1]
if str(SCRIPTS_DIR) not in sys.path:
sys.path.insert(0, str(SCRIPTS_DIR))
import lxc_apps # noqa: E402
class ScheduledUpdateRecordTests(unittest.TestCase):
def setUp(self):
self.temp_dir = tempfile.TemporaryDirectory()
self.addCleanup(self.temp_dir.cleanup)
self.apps_dir_patch = mock.patch.object(
lxc_apps, '_APPS_DIR', self.temp_dir.name,
)
self.apps_dir_patch.start()
self.addCleanup(self.apps_dir_patch.stop)
self.vmid = 9911
self.assertTrue(lxc_apps._write_sidecar(self.vmid, {
'vmid': self.vmid,
'apps': [],
'schedule': {'enabled': True, 'cron': '0 3 * * *'},
}))
def test_record_keeps_log_and_reboot_evidence(self):
self.assertTrue(lxc_apps.record_schedule_run(
self.vmid,
'success',
'both',
log_name='../9911-scheduled-' + ('a' * 32) + '.log',
reboot_required=True,
reboot_packages=['linux-image-amd64', 'libc6'],
))
schedule = lxc_apps._read_sidecar(self.vmid)['schedule']
self.assertEqual(
schedule['last_run_log'],
'9911-scheduled-' + ('a' * 32) + '.log',
)
self.assertTrue(schedule['last_run_reboot_required'])
self.assertEqual(
schedule['last_run_reboot_packages'],
['linux-image-amd64', 'libc6'],
)
def test_lifecycle_clear_preserves_run_and_log(self):
lxc_apps.record_schedule_run(
self.vmid,
'success',
'os',
log_name='9911-scheduled-' + ('b' * 32) + '.log',
reboot_required=True,
reboot_packages=['linux-image-amd64'],
)
self.assertTrue(lxc_apps.clear_schedule_reboot_required(self.vmid))
schedule = lxc_apps._read_sidecar(self.vmid)['schedule']
self.assertFalse(schedule['last_run_reboot_required'])
self.assertNotIn('last_run_reboot_packages', schedule)
self.assertEqual(schedule['last_run_status'], 'success')
self.assertIn('last_run_log', schedule)
if __name__ == '__main__':
unittest.main()
@@ -0,0 +1,160 @@
import sqlite3
import sys
import tempfile
import threading
import time
import unittest
from pathlib import Path
from unittest import mock
SCRIPTS_DIR = Path(__file__).resolve().parents[1]
if str(SCRIPTS_DIR) not in sys.path:
sys.path.insert(0, str(SCRIPTS_DIR))
import notification_manager # noqa: E402
from notification_events import NotificationEvent # noqa: E402
class RecordingChannel:
def __init__(self, results=None):
self.results = list(results or [True])
self.calls = 0
self.lock = threading.Lock()
def send(self, title, body, severity, data):
with self.lock:
self.calls += 1
success = self.results.pop(0) if self.results else True
time.sleep(0.1)
return {'success': success, 'error': '' if success else 'temporary failure'}
class NotificationDeliveryDedupTests(unittest.TestCase):
def setUp(self):
self.temp_dir = tempfile.TemporaryDirectory()
self.addCleanup(self.temp_dir.cleanup)
self.db_path = Path(self.temp_dir.name) / 'health_monitor.db'
conn = sqlite3.connect(str(self.db_path))
conn.execute('''
CREATE TABLE notification_last_sent (
fingerprint TEXT PRIMARY KEY,
last_sent_ts INTEGER NOT NULL,
count INTEGER DEFAULT 1
)
''')
conn.execute('''
CREATE TABLE notification_history (
id INTEGER PRIMARY KEY AUTOINCREMENT,
event_type TEXT NOT NULL,
channel TEXT NOT NULL,
title TEXT,
message TEXT,
severity TEXT,
sent_at TEXT NOT NULL,
success INTEGER DEFAULT 1,
error_message TEXT,
source TEXT DEFAULT 'server'
)
''')
conn.commit()
conn.close()
self.db_patch = mock.patch.object(
notification_manager, 'DB_PATH', self.db_path,
)
self.db_patch.start()
self.addCleanup(self.db_patch.stop)
self.ai_context_patch = mock.patch.object(
notification_manager, 'enrich_context_for_ai', return_value='',
)
self.ai_context_patch.start()
self.addCleanup(self.ai_context_patch.stop)
self.ai_rewrite_patch = mock.patch.object(
notification_manager, '_format_with_ai_bounded', return_value=None,
)
self.ai_rewrite_patch.start()
self.addCleanup(self.ai_rewrite_patch.stop)
def _manager(self, channel):
manager = notification_manager.NotificationManager()
manager._enabled = True
manager._config = {
'email.enabled': 'true',
'email.rich_format': 'false',
'ai_enabled': 'false',
}
manager._channels = {'email': channel}
return manager
@staticmethod
def _event():
return NotificationEvent(
event_type='lxc_update_applied',
severity='INFO',
data={
'hostname': 'pve-test',
'vmid': 210,
'ct_name': 'docker-frontend',
'target': 'Docker Engine',
'result': 'succeeded',
'duration': '10s',
'details': 'Update completed',
},
source='manual',
entity='ct',
entity_id='210:same-run',
)
def test_concurrent_managers_deliver_same_fingerprint_once(self):
channel = RecordingChannel()
managers = [self._manager(channel), self._manager(channel)]
barrier = threading.Barrier(3)
def dispatch(manager):
barrier.wait()
manager._dispatch_event(self._event())
threads = [
threading.Thread(target=dispatch, args=(manager,))
for manager in managers
]
for thread in threads:
thread.start()
barrier.wait()
for thread in threads:
thread.join(timeout=5)
self.assertEqual(channel.calls, 1)
conn = sqlite3.connect(str(self.db_path))
history_count = conn.execute(
'SELECT COUNT(*) FROM notification_history WHERE success = 1'
).fetchone()[0]
claim_count = conn.execute(
'SELECT COUNT(*) FROM notification_delivery_claims'
).fetchone()[0]
conn.close()
self.assertEqual(history_count, 1)
self.assertEqual(claim_count, 0)
def test_failed_delivery_releases_claim_for_retry(self):
channel = RecordingChannel([False, True])
manager = self._manager(channel)
manager._dispatch_event(self._event())
manager._dispatch_event(self._event())
self.assertEqual(channel.calls, 2)
conn = sqlite3.connect(str(self.db_path))
successes = conn.execute(
'SELECT success FROM notification_history ORDER BY id'
).fetchall()
claim_count = conn.execute(
'SELECT COUNT(*) FROM notification_delivery_claims'
).fetchone()[0]
conn.close()
self.assertEqual(successes, [(0,), (1,)])
self.assertEqual(claim_count, 0)
if __name__ == '__main__':
unittest.main()
@@ -0,0 +1,114 @@
import json
import sys
import tempfile
import unittest
from pathlib import Path
from queue import Queue
from unittest import mock
SCRIPTS_DIR = Path(__file__).resolve().parents[1]
if str(SCRIPTS_DIR) not in sys.path:
sys.path.insert(0, str(SCRIPTS_DIR))
import notification_events # noqa: E402
import notification_templates # noqa: E402
import post_install_versions # noqa: E402
class NotificationUpdatePolicyTests(unittest.TestCase):
def test_apt_listchanges_system_mail_has_its_own_update_event(self):
watcher = notification_events.ProxmoxHookWatcher(Queue())
classified = watcher._classify_pve(
"system-mail",
"info",
"Novedades de apt-listchanges para amd",
"zfs-linux recommends that all users update absolute paths",
)
self.assertEqual(classified, ("apt_listchanges", "node", ""))
def test_regular_system_mail_remains_available(self):
watcher = notification_events.ProxmoxHookWatcher(Queue())
with mock.patch.object(notification_events, "_record_smartd_observation_impl"):
classified = watcher._classify_pve(
"system-mail",
"warning",
"SMART error (CurrentPendingSector) detected on host",
"Device: /dev/sda",
)
self.assertEqual(classified, ("system_mail", "node", ""))
def test_apt_listchanges_body_is_preserved_and_attributed(self):
watcher = notification_events.ProxmoxHookWatcher(Queue())
upstream = (
"zfs-linux (2.2.4-2) unstable; urgency=medium\n\n"
" Package-maintainer recommendation.\n\n"
" -- Maintainer <maintainer@example.com>"
)
result = watcher.process_webhook({
"title": "apt-listchanges: News for host",
"message": upstream,
"severity": "info",
"fields": {"type": "system-mail", "hostname": "pve-test"},
})
self.assertTrue(result["accepted"])
event = watcher._queue.get_nowait()
self.assertEqual(event.event_type, "apt_listchanges")
self.assertEqual(event.data["reason"], upstream)
rendered = notification_templates.render_template(
event.event_type,
event.data,
)
self.assertIn("not a ProxMenux recommendation", rendered["body_text"])
self.assertIn(upstream, rendered["body_text"])
def test_smaller_pending_subset_does_not_notify_again(self):
updates = [
{"key": "persistent_network", "available_version": "1.2"},
]
notified = {
"log2ram": {"1.4"},
"persistent_network": {"1.2"},
}
self.assertEqual(
notification_events._new_post_install_update_versions(updates, notified),
{},
)
def test_new_version_of_existing_tool_is_detected(self):
updates = [
{"key": "persistent_network", "available_version": "1.3"},
]
notified = {"persistent_network": {"1.2"}}
self.assertEqual(
notification_events._new_post_install_update_versions(updates, notified),
{"persistent_network": "1.3"},
)
def test_notified_versions_share_the_existing_snapshot_file(self):
with tempfile.TemporaryDirectory() as temporary:
snapshot_path = Path(temporary) / "updates_available.json"
cache = {
"scanned_at": 123.0,
"updates": [
{"key": "log2ram", "available_version": "1.4"},
],
}
with mock.patch.object(post_install_versions, "_UPDATES_JSON", snapshot_path), \
mock.patch.object(post_install_versions, "_cache", cache):
post_install_versions.save_notified_versions(
{"log2ram": {"1.3", "1.4"}}
)
self.assertEqual(
post_install_versions.load_notified_versions(),
{"log2ram": {"1.3", "1.4"}},
)
payload = json.loads(snapshot_path.read_text(encoding="utf-8"))
self.assertEqual(payload["scanned_at"], 123.0)
self.assertEqual(payload["updates"], cache["updates"])
self.assertEqual(payload["notified_versions"]["log2ram"], ["1.3", "1.4"])
if __name__ == "__main__":
unittest.main()
@@ -0,0 +1,37 @@
import sys
from pathlib import Path
import unittest
from unittest.mock import MagicMock, patch
SCRIPTS = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(SCRIPTS))
import lxc_apps
class AdguardSetupTests(unittest.TestCase):
def test_non_adguard_is_not_probed(self):
with patch.object(lxc_apps, '_oci_instance_meta', return_value={'template_id': 'image-frigate'}), \
patch.object(lxc_apps.subprocess, 'run') as run:
self.assertFalse(lxc_apps.oci_adguard_setup_available(190))
run.assert_not_called()
def test_setup_returns_true_only_for_http_response(self):
process = MagicMock(stdout='192.168.0.42\n')
connection = MagicMock()
connection.__enter__.return_value.recv.return_value = b'HTTP/1.1 302 Found\r\n'
with patch.object(lxc_apps, '_oci_instance_meta', return_value={'template_id': 'image-adguard-home'}), \
patch.object(lxc_apps.subprocess, 'run', return_value=process), \
patch.object(lxc_apps.socket, 'create_connection', return_value=connection) as connect:
self.assertTrue(lxc_apps.oci_adguard_setup_available(190))
connect.assert_called_once_with(('192.168.0.42', 3000), timeout=1)
def test_closed_setup_port_is_not_offered(self):
process = MagicMock(stdout='192.168.0.42\n')
with patch.object(lxc_apps, '_oci_instance_meta', return_value={'template_id': 'image-adguard-home'}), \
patch.object(lxc_apps.subprocess, 'run', return_value=process), \
patch.object(lxc_apps.socket, 'create_connection', side_effect=OSError):
self.assertFalse(lxc_apps.oci_adguard_setup_available(190))
if __name__ == '__main__':
unittest.main()
@@ -0,0 +1,55 @@
#!/usr/bin/env python3
"""Regression tests based on LeidenSpain's real LXC OOM block."""
import sys
import unittest
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
from proxmox_known_errors import ( # noqa: E402
analyze_oom_event,
format_oom_diagnosis,
get_error_context,
)
LXC_OOM = """
apt-get invoked oom-killer: gfp_mask=0x101cca, order=3, oom_score_adj=0
memory: usage 131056kB, limit 131072kB, failcnt 1922
swap: usage 0kB, limit 0kB, failcnt 0
Memory cgroup stats for /lxc/108:
oom-kill:constraint=CONSTRAINT_MEMCG,nodemask=(null),cpuset=ns,mems_allowed=0,oom_memcg=/lxc/108,task_memcg=/lxc/108/ns/.lxc,task=apt-get,pid=1183877,uid=100000
Memory cgroup out of memory: Killed process 1183877 (apt-get) total-vm:79700kB, anon-rss:65056kB
"""
class OomDiagnosticsTest(unittest.TestCase):
def test_lxc_memcg_scope_and_limits(self):
result = analyze_oom_event(LXC_OOM)
self.assertIsNotNone(result)
self.assertEqual(result['scope'], 'lxc')
self.assertEqual(result['ctid'], '108')
self.assertEqual(result['constraint'], 'CONSTRAINT_MEMCG')
self.assertEqual(result['memory_usage_kib'], 131056)
self.assertEqual(result['memory_limit_kib'], 131072)
self.assertEqual(result['swap_limit_kib'], 0)
self.assertEqual(result['victim_process'], 'apt-get')
self.assertEqual(result['victim_pid'], '1183877')
def test_diagnosis_does_not_blame_host(self):
diagnosis = format_oom_diagnosis(analyze_oom_event(LXC_OOM))
self.assertIn('LXC 108', diagnosis)
self.assertIn('not a host-wide OOM', diagnosis)
self.assertIn('128.0 MiB used of 128.0 MiB', diagnosis)
self.assertIn('Killed process: apt-get', diagnosis)
def test_known_error_context_includes_event_evidence(self):
context = get_error_context(LXC_OOM, category='memory', detail_level='detailed')
self.assertIn('Event analysis:', context)
self.assertIn('LXC 108 memory cgroup', context)
if __name__ == '__main__':
unittest.main()