mirror of
https://github.com/MacRimi/ProxMenux.git
synced 2026-10-08 22:46:41 +00:00
feat(oci): watchdog that restarts an application when it stops on its own
- New service that starts again an OCI application that stopped on its own, such as Frigate after Save & Restart. A stop or shutdown asked by the user, a backup, a migration and a ProxMenux operation are never undone. - The installer asks it in both modes and proposes what the recipe declares; it is changed from the management menu or with the Watchdog switch of the container modal in the Monitor. - Notifications when the watchdog restarts an application and when it keeps stopping. - The service is recorded in the change journal, and the feature is documented.
This commit is contained in:
@@ -0,0 +1,305 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Start again an OCI application that stopped on its own.
|
||||
|
||||
Proxmox has no restart policy: when the process of an application container
|
||||
ends, the container stays stopped. This service watches the containers whose
|
||||
record asks for it and starts again the one that stopped without anybody
|
||||
asking: a stop, a shutdown, a backup, a migration or a ProxMenux operation is
|
||||
never undone, and neither is a container that was already stopped when the
|
||||
service first saw it.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import json
|
||||
from pathlib import Path
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
|
||||
import oci_instances as instances
|
||||
import oci_operation_notice
|
||||
|
||||
UNIT = 'proxmenux-oci-watchdog.service'
|
||||
UNIT_FILE = Path('/etc/systemd/system') / UNIT
|
||||
JOURNAL = Path('/usr/local/share/proxmenux/scripts/global/pmx_journal.sh')
|
||||
CGROUPS = Path('/sys/fs/cgroup/lxc')
|
||||
CONFIGS = Path('/etc/pve/lxc')
|
||||
HA_RESOURCES = Path('/etc/pve/ha/resources.cfg')
|
||||
TASKS = (Path('/var/log/pve/tasks/active'), Path('/var/log/pve/tasks/index'))
|
||||
INTERVAL = 10
|
||||
# The clock of a task has one second of resolution.
|
||||
TOLERANCE = 2
|
||||
# An application that ran this long before stopping is not in a crash loop:
|
||||
# its restart, such as the one it asks for after saving its settings, is
|
||||
# immediate however often it happens.
|
||||
STABLE = 20
|
||||
# One notice of a restart per application in this time, however often it restarts.
|
||||
NOTICE_EVERY = 3600
|
||||
# Wait before each new attempt; after the last one the application is left stopped.
|
||||
DELAYS = (0, 30, 60, 120, 300)
|
||||
# The marks of an operation are kept for the notices that arrive after it ends.
|
||||
OPERATION_GRACE = 120
|
||||
STOP_TASKS = {'vzstop', 'vzshutdown', 'vzsuspend', 'vzreboot', 'vzdestroy', 'vzmigrate', 'vzrestore', 'vzdump'}
|
||||
HOST_TASKS = {'stopall', 'migrateall'}
|
||||
|
||||
|
||||
def new_state():
|
||||
return {'seen': None, 'since': None, 'attempts': 0, 'retry_at': 0.0, 'settled': False, 'noticed': None}
|
||||
|
||||
|
||||
def notice_due(state, now):
|
||||
"""Whether to tell that the application was restarted: the first time,
|
||||
and then at most once in a while."""
|
||||
if state['noticed'] is not None and now - state['noticed'] < NOTICE_EVERY:
|
||||
return False
|
||||
state['noticed'] = now
|
||||
return True
|
||||
|
||||
|
||||
def parse_tasks(text, active):
|
||||
"""The tasks of a Proxmox task list as (type, id, end); `end` is None
|
||||
while the task runs. The list of active tasks carries one more column."""
|
||||
tasks = []
|
||||
for line in text.splitlines():
|
||||
fields = line.split()
|
||||
if not fields or not fields[0].startswith('UPID:'):
|
||||
continue
|
||||
parts = fields[0].split(':')
|
||||
if len(parts) < 8:
|
||||
continue
|
||||
stamp = fields[2:3] if active else fields[1:2]
|
||||
try:
|
||||
end = int(stamp[0], 16) if stamp else None
|
||||
except ValueError:
|
||||
continue
|
||||
tasks.append((parts[5], parts[6], end))
|
||||
return tasks
|
||||
|
||||
|
||||
def stop_requested(tasks, vmid, seen):
|
||||
"""Whether somebody asked for the container, or for every guest of the
|
||||
host, to stop: the request is still running or ended after the container
|
||||
was last seen running."""
|
||||
return any((kind in HOST_TASKS or (kind in STOP_TASKS and target == str(vmid)))
|
||||
and (end is None or end >= seen - TOLERANCE) for kind, target, end in tasks)
|
||||
|
||||
|
||||
def decide(state, running, now, busy, requested):
|
||||
"""One look at one container. Returns 'start' when it has to be started
|
||||
again, 'give-up' when it keeps stopping, or None. `busy` and `requested`
|
||||
are only called for a container that was running and no longer is."""
|
||||
if running:
|
||||
if state['since'] is None:
|
||||
state['since'] = now
|
||||
if now - state['since'] >= STABLE:
|
||||
state['attempts'] = 0
|
||||
state.update(seen=now, settled=False)
|
||||
return None
|
||||
state['since'] = None
|
||||
if state['seen'] is None or state['settled']:
|
||||
return None
|
||||
if busy():
|
||||
return None
|
||||
if requested(state['seen']):
|
||||
state['settled'] = True
|
||||
return None
|
||||
if now < state['retry_at']:
|
||||
return None
|
||||
if state['attempts'] >= len(DELAYS):
|
||||
state['settled'] = True
|
||||
return 'give-up'
|
||||
state['attempts'] += 1
|
||||
state['retry_at'] = now + (DELAYS[state['attempts']] if state['attempts'] < len(DELAYS) else DELAYS[-1])
|
||||
return 'start'
|
||||
|
||||
|
||||
def watched(root):
|
||||
"""The installed applications that asked for the watchdog: {vmid: name}."""
|
||||
found = {}
|
||||
for path in sorted(root.glob('*/oci-compose.json')):
|
||||
try:
|
||||
record = json.loads(path.read_text())
|
||||
vmid = int(record['vmid'])
|
||||
except (OSError, ValueError, KeyError, TypeError):
|
||||
continue
|
||||
if record.get('status') != 'installed' or record.get('deployment', {}).get('watchdog') is not True:
|
||||
continue
|
||||
if record.get('pending_stack_transaction') or record.get('pending_transaction'):
|
||||
continue
|
||||
title = record.get('template', {}).get('catalog_ui', {}).get('title')
|
||||
found[vmid] = (title.get('en_US') if isinstance(title, dict) else title) or f'CT {vmid}'
|
||||
return found
|
||||
|
||||
|
||||
def is_running(vmid):
|
||||
return (CGROUPS / str(vmid)).is_dir()
|
||||
|
||||
|
||||
def is_busy(vmid, now):
|
||||
"""Something is working on the container or on the host: look again later."""
|
||||
try:
|
||||
config = (CONFIGS / f'{vmid}.conf').read_text(errors='replace')
|
||||
except OSError:
|
||||
return True
|
||||
if any(line.startswith('lock:') for line in config.split('\n[', 1)[0].splitlines()):
|
||||
return True
|
||||
try:
|
||||
mark = json.loads((oci_operation_notice.MARKERS / str(vmid)).read_text())
|
||||
if mark.get('ended') is None or now - float(mark['ended']) < OPERATION_GRACE:
|
||||
return True
|
||||
except (OSError, ValueError, TypeError, KeyError):
|
||||
pass
|
||||
try:
|
||||
if any(line.split(':', 1)[0].strip() == 'ct' and line.split(':', 1)[1].strip() == str(vmid)
|
||||
for line in HA_RESOURCES.read_text().splitlines() if ':' in line):
|
||||
return True
|
||||
except OSError:
|
||||
pass
|
||||
state = subprocess.run(['systemctl', 'is-system-running'], capture_output=True, text=True, check=False)
|
||||
return state.stdout.strip() == 'stopping'
|
||||
|
||||
|
||||
def read_tasks():
|
||||
tasks = []
|
||||
for path in TASKS:
|
||||
try:
|
||||
with path.open('rb') as handle:
|
||||
handle.seek(0, 2)
|
||||
handle.seek(max(0, handle.tell() - 65536))
|
||||
tasks += parse_tasks(handle.read().decode(errors='replace'), path.name == 'active')
|
||||
except OSError:
|
||||
continue
|
||||
return tasks
|
||||
|
||||
|
||||
def start(vmid):
|
||||
# In a scope of its own: what Proxmox leaves running for the container,
|
||||
# such as its DHCP client, must not end when this service is restarted.
|
||||
result = subprocess.run(['systemd-run', '--scope', '--quiet', '--collect', 'pct', 'start', str(vmid)],
|
||||
capture_output=True, text=True, check=False, timeout=300)
|
||||
return result.returncode == 0
|
||||
|
||||
|
||||
def look(states, now):
|
||||
apps = watched(instances.ROOT)
|
||||
for vmid in set(states) - set(apps):
|
||||
del states[vmid]
|
||||
for vmid, name in apps.items():
|
||||
state = states.setdefault(vmid, new_state())
|
||||
action = decide(state, is_running(vmid), now, lambda: is_busy(vmid, now),
|
||||
lambda seen: stop_requested(read_tasks(), vmid, seen))
|
||||
data = {'app_name': name, 'vmid': vmid, 'containers': f'CT {vmid}'}
|
||||
if action == 'start':
|
||||
print(f'CT {vmid} ({name}) stopped on its own; starting it again (attempt {state["attempts"]})', flush=True)
|
||||
try:
|
||||
started = start(vmid)
|
||||
except (OSError, subprocess.SubprocessError):
|
||||
started = False
|
||||
if started and notice_due(state, now):
|
||||
oci_operation_notice.notify('oci_watchdog_restarted', data)
|
||||
elif action == 'give-up':
|
||||
print(f'CT {vmid} ({name}) keeps stopping; it is left stopped', flush=True)
|
||||
oci_operation_notice.notify('oci_watchdog_failed', data)
|
||||
|
||||
|
||||
def run():
|
||||
states = {}
|
||||
while True:
|
||||
try:
|
||||
look(states, time.time())
|
||||
except (OSError, ValueError) as error:
|
||||
print(f'watchdog: {error}', file=sys.stderr, flush=True)
|
||||
time.sleep(INTERVAL)
|
||||
|
||||
|
||||
def _journaled(unit):
|
||||
"""Write and enable the unit through the change journal of ProxMenux, so
|
||||
the Changes tab of the Monitor shows it. False when the journal is not
|
||||
installed on this host."""
|
||||
if not JOURNAL.is_file():
|
||||
return False
|
||||
script = ('source "$1" && pmx_journal_context "oci_watchdog" "1.0" "oci_watchdog.py" '
|
||||
'&& pmx_write_file "$2" && systemctl daemon-reload && pmx_enable_service "$3"')
|
||||
result = subprocess.run(['bash', '-c', script, 'bash', str(JOURNAL), str(UNIT_FILE), UNIT],
|
||||
input=unit, text=True, capture_output=True, check=False)
|
||||
return result.returncode == 0
|
||||
|
||||
|
||||
def ensure_service():
|
||||
"""Install the service and leave it running. A running one is only
|
||||
restarted when its unit changed; the ProxMenux installer restarts it when
|
||||
it replaces this program."""
|
||||
unit = ('[Unit]\n'
|
||||
'Description=ProxMenux OCI watchdog\n'
|
||||
'After=pve-guests.service\n\n'
|
||||
'[Service]\n'
|
||||
'Type=simple\n'
|
||||
f'ExecStart=/usr/bin/python3 {Path(__file__).resolve()} run\n'
|
||||
'Restart=on-failure\n'
|
||||
'RestartSec=30\n\n'
|
||||
'[Install]\n'
|
||||
'WantedBy=multi-user.target\n')
|
||||
changed = not UNIT_FILE.is_file() or UNIT_FILE.read_text() != unit
|
||||
if not _journaled(unit):
|
||||
if changed:
|
||||
UNIT_FILE.write_text(unit)
|
||||
subprocess.run(['systemctl', 'daemon-reload'], check=False)
|
||||
subprocess.run(['systemctl', 'enable', UNIT], check=False, capture_output=True)
|
||||
subprocess.run(['systemctl', 'restart' if changed else 'start', UNIT], check=False, capture_output=True)
|
||||
|
||||
|
||||
def application(root, vmid):
|
||||
"""The containers of the application `vmid` belongs to: itself, or every
|
||||
member of its stack."""
|
||||
record = instances.read(root, vmid)
|
||||
primary_id = int((record.get('stack_member') or {}).get('primary_vmid') or vmid)
|
||||
primary = record if primary_id == vmid else instances.read(root, primary_id)
|
||||
members = [int(member['vmid']) for member in (primary.get('stack') or {}).get('members') or []]
|
||||
return members or [vmid]
|
||||
|
||||
|
||||
def set_watchdog(root, vmid, enabled):
|
||||
"""Turn the watchdog of an application on or off. The choice is kept in
|
||||
the record of each of its containers, so updates and backups carry it."""
|
||||
import oci_carried_record
|
||||
with instances.locked(root):
|
||||
vmids = application(root, vmid)
|
||||
for member in vmids:
|
||||
record = instances.read(root, member)
|
||||
record.setdefault('deployment', {})['watchdog'] = bool(enabled)
|
||||
instances.write(instances.location(root, member), record)
|
||||
oci_carried_record.carry(root, member, mount_stopped=False)
|
||||
if enabled:
|
||||
ensure_service()
|
||||
return vmids
|
||||
|
||||
|
||||
def main():
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
commands = parser.add_subparsers(dest='command', required=True)
|
||||
commands.add_parser('run')
|
||||
commands.add_parser('install')
|
||||
change = commands.add_parser('set')
|
||||
change.add_argument('vmid', type=int)
|
||||
change.add_argument('state', choices=['on', 'off'])
|
||||
args = parser.parse_args()
|
||||
if args.command == 'install':
|
||||
ensure_service()
|
||||
elif args.command == 'set':
|
||||
try:
|
||||
vmids = set_watchdog(instances.ROOT, args.vmid, args.state == 'on')
|
||||
except BlockingIOError:
|
||||
print('Another OCI operation is using the instance registry.', file=sys.stderr)
|
||||
return 3
|
||||
except (OSError, ValueError, KeyError) as error:
|
||||
print(str(error) or type(error).__name__, file=sys.stderr)
|
||||
return 1
|
||||
print(json.dumps({'vmids': vmids, 'watchdog': args.state == 'on'}))
|
||||
else:
|
||||
run()
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
raise SystemExit(main())
|
||||
@@ -346,6 +346,8 @@ def _deployment_summary_text(template: dict[str, Any], deployment: dict[str, Any
|
||||
lines.append("")
|
||||
row(translate("Start"), f"{translate('when finished')}: {_yes_no(plan.get('start_after_create'))} · "
|
||||
f"{translate('with Proxmox')}: {_yes_no(plan.get('onboot'))}")
|
||||
if "watchdog" in plan:
|
||||
row(translate("Watchdog"), _yes_no(plan["watchdog"]))
|
||||
return "\n".join(lines)
|
||||
|
||||
|
||||
@@ -383,6 +385,14 @@ def _print_installation_summary(result: dict[str, Any], images_removed: str | No
|
||||
|
||||
# ---------------------------------------------------------------- menus
|
||||
|
||||
def watchdog_default(template: dict[str, Any]) -> bool:
|
||||
"""Whether the recipe of the application asks Docker to restart it."""
|
||||
policies = [template.get("container_contract", {}).get("restart")]
|
||||
policies += [service.get("compose", {}).get("restart")
|
||||
for service in template.get("compose_stack", {}).get("services", []) if isinstance(service, dict)]
|
||||
return any(policy in ("always", "unless-stopped", "on-failure") for policy in policies)
|
||||
|
||||
|
||||
def _install(catalog: Catalog, ui, item: dict[str, Any], mode: str) -> None:
|
||||
install_template(ui, catalog.compose(item["id"]), item["id"], mode)
|
||||
|
||||
@@ -397,6 +407,10 @@ def install_template(ui, template: dict[str, Any], identifier: str, mode: str) -
|
||||
candidate = copy.deepcopy(template)
|
||||
try:
|
||||
deployment = build_deployment(candidate, wizard, mode)
|
||||
# Every installation decides it, the default one as well.
|
||||
deployment["watchdog"] = wizard.confirm(
|
||||
translate("Put this application under watchdog? It is restarted automatically when it crashes."),
|
||||
watchdog_default(candidate))
|
||||
approved = wizard.review(_deployment_summary_text(candidate, deployment),
|
||||
translate("Installation summary"),
|
||||
question=translate("Install with this configuration?"))
|
||||
@@ -413,6 +427,8 @@ def install_template(ui, template: dict[str, Any], identifier: str, mode: str) -
|
||||
if not approved:
|
||||
return None
|
||||
template = candidate
|
||||
# The engine installs the application; the watchdog is turned on once it is there.
|
||||
watchdog = bool(deployment.pop("watchdog", False))
|
||||
console.show_logo()
|
||||
console.msg_title(f"{source_text(template['catalog_ui']['title']) or identifier} · {APP_TITLE}")
|
||||
try:
|
||||
@@ -427,6 +443,12 @@ def install_template(ui, template: dict[str, Any], identifier: str, mode: str) -
|
||||
vmids = {int(v) for v in [result.get("vmid"), *(result.get("stack_vmids") or {}).values()] if v}
|
||||
_, removed = images.offer_removal(ui, sorted(vmids))
|
||||
_print_installation_summary(result, removed)
|
||||
if watchdog:
|
||||
from .management import set_watchdog
|
||||
if set_watchdog(PROJECT_ROOT, sorted(vmids), True):
|
||||
console.msg_ok(translate("Watchdog enabled: the application is restarted when it crashes."))
|
||||
else:
|
||||
console.msg_warn(translate("The watchdog could not be enabled. Turn it on from Manage installed OCI applications."))
|
||||
console.wait_for_enter(translate("Press Enter to return to the menu..."))
|
||||
return result
|
||||
|
||||
|
||||
@@ -269,6 +269,39 @@ def _interactive_management(project, ui):
|
||||
manage_instance(project, ui, row)
|
||||
|
||||
|
||||
def set_watchdog(project, vmids, enabled):
|
||||
"""Turn the watchdog of the applications these containers belong to on or
|
||||
off. False when the registry is busy or a record cannot be read."""
|
||||
sys.path.insert(0, str(project / 'remote'))
|
||||
import oci_instances as instances
|
||||
import oci_watchdog
|
||||
done = set()
|
||||
for vmid in vmids:
|
||||
if vmid in done:
|
||||
continue
|
||||
try:
|
||||
done.update(oci_watchdog.set_watchdog(instances.ROOT, vmid, enabled))
|
||||
except (OSError, ValueError, KeyError):
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
def _toggle_watchdog(project, ui, vmid, enabled):
|
||||
"""Ask and change the watchdog of an application from its menu."""
|
||||
question = (translate('This application is under watchdog: it is restarted automatically when it crashes. Turn the watchdog off?')
|
||||
if enabled else
|
||||
translate('Put this application under watchdog? It is restarted automatically when it crashes. A stop or a shutdown you ask for is never undone.'))
|
||||
if not ui.confirm(question, not enabled):
|
||||
return False
|
||||
if not set_watchdog(project, [vmid], not enabled):
|
||||
ui.message(translate('Another OCI operation is using the instance registry. Wait for it to finish and open this menu again; no container is modified.'),
|
||||
translate('OCI management'))
|
||||
return False
|
||||
ui.message(translate('Watchdog disabled.') if enabled else translate('Watchdog enabled: the application is restarted when it crashes.'),
|
||||
translate('OCI management'))
|
||||
return True
|
||||
|
||||
|
||||
def manage_instance(project, ui, row, action=None, lifecycle_args=()):
|
||||
"""What the menu does with one instance once it is selected. `action`
|
||||
skips the choice of operation, as ProxMenux Monitor does; the extra
|
||||
@@ -292,15 +325,21 @@ def manage_instance(project, ui, row, action=None, lifecycle_args=()):
|
||||
if row['status'] != 'installed' or row['reason'] != 'matched':
|
||||
ui.message(translate('The instance identity or status must be reviewed before updating.'), translate('OCI management'))
|
||||
return False
|
||||
if action is None:
|
||||
action = ui.choose(translate('Manage OCI'), [('update', translate('Update the image with the saved configuration')),
|
||||
('modify', translate('Modify: edit resources, network, paths and GPU')),
|
||||
('remove', translate('Remove: delete the application and its containers'))], 'update')
|
||||
if action is None:
|
||||
return False
|
||||
sys.path.insert(0, str(project / 'remote'))
|
||||
import oci_instances as instances
|
||||
record = instances.read(instances.ROOT, row['vmid'])
|
||||
watched = record.get('deployment', {}).get('watchdog') is True
|
||||
watchdog_label = (translate('Watchdog (on): restart the application when it crashes') if watched
|
||||
else translate('Watchdog (off): restart the application when it crashes'))
|
||||
if action is None:
|
||||
action = ui.choose(translate('Manage OCI'), [('update', translate('Update the image with the saved configuration')),
|
||||
('modify', translate('Modify: edit resources, network, paths and GPU')),
|
||||
('watchdog', watchdog_label),
|
||||
('remove', translate('Remove: delete the application and its containers'))], 'update')
|
||||
if action is None:
|
||||
return False
|
||||
if action == 'watchdog':
|
||||
return _toggle_watchdog(project, ui, row['vmid'], watched)
|
||||
if action == 'remove':
|
||||
return _remove(project, ui, row['vmid'])
|
||||
import oci_instance_reconcile as reconcile
|
||||
@@ -489,10 +528,15 @@ def _manage_stack(project, ui, row, action=None, lifecycle_args=()):
|
||||
options.append(('modify', translate('Modify extra paths and devices')))
|
||||
if updatable:
|
||||
options.append(('recreate', translate('Recreate every container with its saved configuration')))
|
||||
watched = primary.get('deployment', {}).get('watchdog') is True
|
||||
options.append(('watchdog', translate('Watchdog (on): restart the application when it crashes') if watched
|
||||
else translate('Watchdog (off): restart the application when it crashes')))
|
||||
options.append(('remove', translate('Remove: delete the application and its containers')))
|
||||
action = ui.choose(translate('Manage OCI stack'), options, options[0][0])
|
||||
if action is None:
|
||||
return False
|
||||
if action == 'watchdog':
|
||||
return _toggle_watchdog(project, ui, primary_id, primary.get('deployment', {}).get('watchdog') is True)
|
||||
if action == 'remove':
|
||||
return _remove(project, ui, primary_id)
|
||||
if action == 'modify':
|
||||
|
||||
@@ -53,6 +53,7 @@ class StackRecreateV1Tests(unittest.TestCase):
|
||||
"base_config_sha256": "saved-config"}}
|
||||
with patch.object(adapter, "validate"), \
|
||||
patch.object(adapter, "state", return_value={"backups": {"138": {"archive": "/backup"}}}), \
|
||||
patch.object(native, "translate", side_effect=lambda text: text), \
|
||||
patch.object(native.member_tx, "apply") as apply:
|
||||
adapter.replace(138, prepared, "transaction-id")
|
||||
|
||||
|
||||
@@ -0,0 +1,209 @@
|
||||
"""The watchdog starts again what stopped on its own, and nothing else."""
|
||||
|
||||
import json
|
||||
from pathlib import Path
|
||||
import sys
|
||||
import tempfile
|
||||
import unittest
|
||||
|
||||
ROOT = Path(__file__).resolve().parents[1]
|
||||
sys.path.insert(0, str(ROOT / "remote"))
|
||||
|
||||
import oci_watchdog as watchdog
|
||||
|
||||
ACTIVE = """UPID:amd:001C34BC:01A809CF:6AC55103:vzshutdown:109:root@pam: 1 6AC55105 OK
|
||||
UPID:amd:001C3370:01A80568:6AC550F8:vzstart:109:root@pam: 1 6AC550FA OK
|
||||
UPID:amd:001C32AC:01A80211:6AC550EF:vzstop:108:root@pam: 1 6AC550F1 unable to stop: timeout
|
||||
UPID:amd:001C3999:01A80999:6AC55200:vzdump:110:root@pam: 0
|
||||
"""
|
||||
INDEX = """UPID:amd:001C32AC:01A80211:6AC550EF:vzstop:109:root@pam: 6AC550F1 OK
|
||||
UPID:amd:001C3370:01A80568:6AC550F8:stopall::root@pam: 6AC55300 unable to stop every guest
|
||||
"""
|
||||
|
||||
|
||||
def never():
|
||||
raise AssertionError("not expected to be asked")
|
||||
|
||||
|
||||
class Watch:
|
||||
"""One container seen through the decisions of the watchdog."""
|
||||
def __init__(self):
|
||||
self.state = watchdog.new_state()
|
||||
|
||||
def look(self, running, now, busy=False, requested=False):
|
||||
return watchdog.decide(self.state, running, now, lambda: busy, lambda seen: requested)
|
||||
|
||||
|
||||
class TaskList(unittest.TestCase):
|
||||
def test_both_lists_are_read_with_their_own_columns(self):
|
||||
self.assertEqual(watchdog.parse_tasks(ACTIVE, True), [
|
||||
("vzshutdown", "109", 0x6AC55105), ("vzstart", "109", 0x6AC550FA),
|
||||
("vzstop", "108", 0x6AC550F1), ("vzdump", "110", None)])
|
||||
self.assertEqual(watchdog.parse_tasks(INDEX, False), [("vzstop", "109", 0x6AC550F1), ("stopall", "", 0x6AC55300)])
|
||||
|
||||
def test_a_request_counts_when_it_ended_after_the_container_was_last_seen_running(self):
|
||||
tasks = watchdog.parse_tasks(ACTIVE, True)
|
||||
self.assertTrue(watchdog.stop_requested(tasks, 109, 0x6AC55104))
|
||||
self.assertFalse(watchdog.stop_requested(tasks, 109, 0x6AC55110))
|
||||
self.assertFalse(watchdog.stop_requested(tasks, 111, 0x6AC55000))
|
||||
|
||||
def test_a_running_request_and_a_stop_of_every_guest_count(self):
|
||||
self.assertTrue(watchdog.stop_requested(watchdog.parse_tasks(ACTIVE, True), 110, 0x6AC55900))
|
||||
self.assertTrue(watchdog.stop_requested(watchdog.parse_tasks(INDEX, False), 555, 0x6AC552F0))
|
||||
|
||||
def test_starting_a_container_is_not_a_request_to_stop_it(self):
|
||||
self.assertFalse(watchdog.stop_requested([("vzstart", "109", None)], 109, 0))
|
||||
|
||||
|
||||
class Decisions(unittest.TestCase):
|
||||
def test_a_container_never_seen_running_is_left_alone(self):
|
||||
watch = Watch()
|
||||
self.assertIsNone(watchdog.decide(watch.state, False, 100, never, never))
|
||||
self.assertIsNone(watchdog.decide(watch.state, False, 5000, never, never))
|
||||
|
||||
def test_a_container_that_stops_on_its_own_is_started_again(self):
|
||||
watch = Watch()
|
||||
watch.look(True, 100)
|
||||
self.assertEqual(watch.look(False, 110), "start")
|
||||
self.assertEqual(watch.state["attempts"], 1)
|
||||
|
||||
def test_a_requested_stop_is_never_undone(self):
|
||||
watch = Watch()
|
||||
watch.look(True, 100)
|
||||
self.assertIsNone(watch.look(False, 110, requested=True))
|
||||
self.assertIsNone(watchdog.decide(watch.state, False, 9000, never, never))
|
||||
watch.look(True, 9100)
|
||||
self.assertEqual(watch.look(False, 9110), "start")
|
||||
|
||||
def test_nothing_is_decided_while_something_works_on_the_container(self):
|
||||
watch = Watch()
|
||||
watch.look(True, 100)
|
||||
self.assertIsNone(watchdog.decide(watch.state, False, 110, lambda: True, never))
|
||||
self.assertIsNone(watch.look(False, 120, requested=True))
|
||||
watch.look(True, 300)
|
||||
self.assertIsNone(watchdog.decide(watch.state, False, 310, lambda: True, never))
|
||||
self.assertEqual(watch.look(False, 320), "start")
|
||||
|
||||
def test_a_stop_asked_after_a_restart_is_respected(self):
|
||||
watch = Watch()
|
||||
watch.look(True, 100)
|
||||
self.assertEqual(watch.look(False, 110), "start")
|
||||
watch.look(True, 120)
|
||||
self.assertIsNone(watch.look(False, 130, requested=True))
|
||||
self.assertIsNone(watchdog.decide(watch.state, False, 5000, never, never))
|
||||
|
||||
def test_an_application_that_keeps_stopping_is_tried_with_longer_waits_and_then_left(self):
|
||||
watch, now, actions = Watch(), 100, []
|
||||
watch.look(True, now)
|
||||
for _ in range(200):
|
||||
now += 10
|
||||
action = watch.look(False, now)
|
||||
if action:
|
||||
actions.append((action, now))
|
||||
self.assertEqual([name for name, _ in actions], ["start"] * 5 + ["give-up"])
|
||||
waits = [later - earlier for (_, earlier), (_, later) in zip(actions, actions[1:])]
|
||||
self.assertEqual(waits, [30, 60, 120, 300, 300])
|
||||
self.assertIsNone(watchdog.decide(watch.state, False, now + 9000, never, never))
|
||||
|
||||
def test_an_application_that_restarts_itself_now_and_then_is_always_started_at_once(self):
|
||||
watch, now = Watch(), 100
|
||||
watch.look(True, now)
|
||||
for _ in range(12):
|
||||
now += 10
|
||||
self.assertEqual(watch.look(False, now), "start")
|
||||
for _ in range(4):
|
||||
now += 10
|
||||
watch.look(True, now)
|
||||
self.assertEqual(watch.state["attempts"], 0)
|
||||
|
||||
def test_a_restart_is_told_once_in_a_while(self):
|
||||
state = watchdog.new_state()
|
||||
self.assertTrue(watchdog.notice_due(state, 1000))
|
||||
self.assertFalse(watchdog.notice_due(state, 1000 + watchdog.NOTICE_EVERY - 1))
|
||||
self.assertTrue(watchdog.notice_due(state, 1000 + watchdog.NOTICE_EVERY))
|
||||
|
||||
def test_running_for_a_while_forgets_the_earlier_attempts(self):
|
||||
watch = Watch()
|
||||
watch.look(True, 100)
|
||||
self.assertEqual(watch.look(False, 110), "start")
|
||||
watch.look(True, 120)
|
||||
watch.look(True, 120 + watchdog.STABLE)
|
||||
self.assertEqual(watch.state["attempts"], 0)
|
||||
self.assertEqual(watch.look(False, 400), "start")
|
||||
self.assertEqual(watch.state["attempts"], 1)
|
||||
|
||||
|
||||
class Registry(unittest.TestCase):
|
||||
def test_only_installed_applications_that_asked_for_it_are_watched(self):
|
||||
with tempfile.TemporaryDirectory() as folder:
|
||||
root = Path(folder)
|
||||
records = {
|
||||
101: {"status": "installed", "deployment": {"watchdog": True},
|
||||
"template": {"catalog_ui": {"title": {"en_US": "Jellyfin"}}}},
|
||||
102: {"status": "installed", "deployment": {"watchdog": False}},
|
||||
103: {"status": "installed", "deployment": {}},
|
||||
104: {"status": "updating", "deployment": {"watchdog": True}},
|
||||
105: {"status": "installed", "deployment": {"watchdog": True}, "pending_stack_transaction": "/journal"},
|
||||
106: {"status": "installed", "deployment": {"watchdog": True}},
|
||||
}
|
||||
for vmid, record in records.items():
|
||||
(root / str(vmid)).mkdir()
|
||||
(root / str(vmid) / "oci-compose.json").write_text(json.dumps({"vmid": vmid, **record}))
|
||||
(root / "107").mkdir()
|
||||
(root / "107" / "oci-compose.json").write_text("{broken")
|
||||
self.assertEqual(watchdog.watched(root), {101: "Jellyfin", 106: "CT 106"})
|
||||
|
||||
|
||||
class Start(unittest.TestCase):
|
||||
def test_the_container_is_started_outside_the_service(self):
|
||||
from unittest.mock import patch
|
||||
with patch.object(watchdog.subprocess, "run") as run:
|
||||
run.return_value.returncode = 0
|
||||
self.assertTrue(watchdog.start(109))
|
||||
run.return_value.returncode = 255
|
||||
self.assertFalse(watchdog.start(109))
|
||||
self.assertEqual(run.call_args.args[0], ["systemd-run", "--scope", "--quiet", "--collect", "pct", "start", "109"])
|
||||
|
||||
|
||||
class Switch(unittest.TestCase):
|
||||
"""Turning the watchdog on or off for a whole application."""
|
||||
def registry(self, folder):
|
||||
root = Path(folder)
|
||||
ids = {101: "11111111-1111-4111-8111-111111111111", 102: "22222222-2222-4222-8222-222222222222",
|
||||
103: "33333333-3333-4333-8333-333333333333"}
|
||||
stack = {"members": [{"vmid": 101}, {"vmid": 102}]}
|
||||
records = {101: {"stack": stack, "stack_member": {"primary_vmid": 101}},
|
||||
102: {"stack_member": {"primary_vmid": 101}}, 103: {}}
|
||||
for vmid, extra in records.items():
|
||||
(root / str(vmid)).mkdir()
|
||||
(root / str(vmid) / "oci-compose.json").write_text(json.dumps({
|
||||
"schema_version": 1, "vmid": vmid, "installation_id": ids[vmid], "status": "installed",
|
||||
"deployment": {"onboot": True}, **extra}))
|
||||
return root
|
||||
|
||||
def flags(self, root):
|
||||
return {int(path.parent.name): json.loads(path.read_text())["deployment"].get("watchdog")
|
||||
for path in root.glob("*/oci-compose.json")}
|
||||
|
||||
def test_every_container_of_the_application_takes_the_choice_and_keeps_the_rest_of_its_record(self):
|
||||
from unittest.mock import patch
|
||||
import oci_carried_record
|
||||
with tempfile.TemporaryDirectory() as folder:
|
||||
root = self.registry(folder)
|
||||
with patch.object(watchdog, "ensure_service") as service, \
|
||||
patch.object(oci_carried_record, "carry", return_value="carried") as carry:
|
||||
self.assertEqual(watchdog.set_watchdog(root, 102, True), [101, 102])
|
||||
self.assertEqual(self.flags(root), {101: True, 102: True, 103: None})
|
||||
service.assert_called_once_with()
|
||||
self.assertEqual([call.args[1] for call in carry.call_args_list], [101, 102])
|
||||
self.assertEqual(watchdog.set_watchdog(root, 103, True), [103])
|
||||
self.assertEqual(watchdog.set_watchdog(root, 101, False), [101, 102])
|
||||
self.assertEqual(self.flags(root), {101: False, 102: False, 103: True})
|
||||
self.assertEqual(service.call_count, 2)
|
||||
record = json.loads((root / "101" / "oci-compose.json").read_text())
|
||||
self.assertEqual(record["deployment"]["onboot"], True)
|
||||
self.assertEqual(record["stack"]["members"], [{"vmid": 101}, {"vmid": 102}])
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
Reference in New Issue
Block a user