#!/usr/bin/env python3 """Durable workstation deployment controller. Install with setup-workstation.sh.""" import argparse import contextlib import fcntl import importlib.util import json import math import os import re import shutil import subprocess import sys import time from pathlib import Path STATE = Path(os.environ.get('HOMELAB_STATE', Path.home() / '.local/state/homelab-deploy')) CONFIG_REPO = Path(os.environ.get('HOMELAB_REPO', '/srv/homelab')) RUN_ID = re.compile(r'[0-9]+-[0-9]+') def command(*args, **kwargs): return subprocess.check_output(args, text=True, **kwargs).strip() # noqa: S603, S607 def atomic_json(path, data): temporary = path.with_suffix('.tmp') temporary.write_text(json.dumps(data, indent=2) + '\n') temporary.chmod(0o600) temporary.replace(path) @contextlib.contextmanager def lock(name): STATE.mkdir(mode=0o700, parents=True, exist_ok=True) with (STATE / name).open('a') as stream: fcntl.flock(stream, fcntl.LOCK_EX) yield def load_module(name, path): spec = importlib.util.spec_from_file_location(name, path) module = importlib.util.module_from_spec(spec) spec.loader.exec_module(module) return module def run_directory(run_id): if not RUN_ID.fullmatch(run_id): raise ValueError('Run ID must be numeric workflow-id and attempt') return STATE / 'runs' / run_id def start(run_id): payload = sys.stdin.buffer.read(256 * 1024 + 1) if len(payload) > 256 * 1024: raise ValueError('Deploy request exceeds 256 KiB') request = json.loads(payload) sha = request['release']['sha'] if not re.fullmatch(r'[0-9a-f]{40}', sha) or request['mode'] not in ('changed', 'full', 'plan'): raise ValueError('Invalid deploy SHA or mode') if not isinstance(request['refresh_images'], bool): raise ValueError('refresh_images must be boolean') directory = run_directory(run_id) with lock('prepare.lock'): if (directory / 'request.json').exists(): if json.loads((directory / 'request.json').read_text()) != request: raise ValueError('Run ID already belongs to a different request') else: directory.mkdir(mode=0o700, parents=True, exist_ok=True) command('git', '-C', str(CONFIG_REPO), 'fetch', '--quiet', 'origin', 'main') command('git', '-C', str(CONFIG_REPO), 'merge-base', '--is-ancestor', sha, 'origin/main') if not (directory / 'source').exists(): command('git', '-C', str(CONFIG_REPO), 'worktree', 'add', '--detach', str(directory / 'source'), sha) if command('git', '-C', str(directory / 'source'), 'rev-parse', 'HEAD') != sha: raise ValueError('Prepared source does not match deploy SHA') release_module = load_module('release', directory / 'source/.gitea/workflows/release.py') release_module.validate_release(request['release'], sha) atomic_json(directory / 'release.json', request['release']) atomic_json(directory / 'request.json', request) if not (directory / 'status.json').exists(): atomic_json(directory / 'status.json', {'state': 'queued', 'stages': {}}) # Starting an existing active or finished ID is idempotent; never re-apply it. if json.loads((directory / 'status.json').read_text())['state'] == 'queued': command('systemctl', '--user', 'start', '--no-block', f'homelab-deploy@{run_id}.service') print(f'Accepted deploy {run_id} ({sha})') def environment(directory): request = json.loads((directory / 'request.json').read_text()) return { **os.environ, 'REPO': str(directory / 'source'), 'CONFIG_REPO': str(CONFIG_REPO), 'RUN_DIR': str(directory), 'DEPLOY_SHA': request['release']['sha'], 'RELEASE_FILE': str(directory / 'release.json'), 'DEPLOY_PLAN': str(directory / 'plan.json'), 'DEPLOY_SNAPSHOT_DIR': str(directory / 'snapshot'), 'REFRESH_IMAGES': str(request['refresh_images']).lower(), 'ROLLOUT_PARALLELISM': '4', } def stage(directory, name, budget): status = json.loads((directory / 'status.json').read_text()) if name in status['stages'] and status['stages'][name].get('result') in ('success', 'failure'): return status['stages'][name]['result'] == 'success' started = time.time() status['stages'][name] = {'result': 'running', 'started': started} atomic_json(directory / 'status.json', status) script = directory / 'source/.gitea/workflows/deploy-stage.sh' with (directory / f'{name}.log').open('a') as log: # timeout kills the whole stage process group, including children, before recovery. result = subprocess.run( # noqa: S603, S607 [ shutil.which('timeout') or '/usr/bin/timeout', '--signal=TERM', '--kill-after=30s', str(budget), 'bash', str(script), name, ], env=environment(directory), stdout=log, stderr=subprocess.STDOUT, check=False, ).returncode status = json.loads((directory / 'status.json').read_text()) status['stages'][name].update( result='success' if result == 0 else 'failure', exit_code=result, seconds=round(time.time() - started) ) atomic_json(directory / 'status.json', status) return result == 0 def make_plan(directory): source = directory / 'source' planner = load_module('deploy_plan', source / '.gitea/workflows/deploy-plan.py') request = json.loads((directory / 'request.json').read_text()) previous = json.loads((STATE / 'last-success.json').read_text()) if (STATE / 'last-success.json').exists() else None # Helm 4 lists every release status by default and removed the --all flag. helm = json.loads(command('helm', 'list', '-A', '-o', 'json')) plan = planner.make_plan(source, CONFIG_REPO, request['release'], previous, request['mode'], helm) if request['refresh_images']: plan['selected']['compose'] = plan['active']['compose'] atomic_json(directory / 'plan.json', plan) if previous: atomic_json(directory / 'previous.json', previous) # Local config is deliberately separate from the immutable Git source. return plan def finish_success(directory, plan): # Repeating finalization after a crash is safe while holding deploy.lock. plan['run_id'] = directory.name path = directory / 'compose-images.json' previous = directory / 'previous.json' plan['compose-images'] = ( json.loads(path.read_text()) if path.exists() else json.loads(previous.read_text()).get('compose-images', {}) if previous.exists() else {} ) configs = STATE / 'compose-configs' configs.mkdir(mode=0o700, exist_ok=True) for config in (directory / 'compose').glob('*.json'): atomic_json(configs / config.name, json.loads(config.read_text())) atomic_json(STATE / 'last-success.json', plan) status = json.loads((directory / 'status.json').read_text()) status['state'] = 'success' atomic_json(directory / 'status.json', status) try: retain_completed(directory) except (OSError, subprocess.CalledProcessError) as error: print(f'Retention deferred: {error}', flush=True) def recover(directory, retry=False): status = json.loads((directory / 'status.json').read_text()) if status['state'] in ('success', 'planned'): return completed = ('doctor', 'validate', 'apply-k8s', 'apply-compose', 'verify-k8s', 'smoke') if all(status['stages'].get(name, {}).get('result') == 'success' for name in completed): finish_success(directory, json.loads((directory / 'plan.json').read_text())) return if retry: for name in ('verify-k8s', 'smoke'): if status['stages'].get(name, {}).get('result') == 'failure': del status['stages'][name] atomic_json(directory / 'status.json', status) snapshot = directory / 'snapshot/current' if snapshot.exists(): stage(directory, 'verify-k8s', 7200) stage(directory, 'smoke', 600) status = json.loads((directory / 'status.json').read_text()) status['state'] = 'failure' atomic_json(directory / 'status.json', status) def execute(run_id): directory = run_directory(run_id) with lock('deploy.lock'): status = json.loads((directory / 'status.json').read_text()) if status['state'] != 'queued': return # A crashed predecessor must be recovered before another apply begins. for other in (STATE / 'runs').iterdir(): if ( other != directory and (other / 'status.json').exists() and json.loads((other / 'status.json').read_text())['state'] == 'running' ): raise ValueError(f'Interrupted deploy {other.name}; run recover first') status['state'] = 'running' atomic_json(directory / 'status.json', status) phase = 'plan' try: plan = make_plan(directory) print( json.dumps({'selected': plan['selected'], 'helm': plan['helm'], 'manual_removals': plan['removed']}), flush=True, ) phase = 'doctor' if not stage(directory, 'doctor', 600): raise RuntimeError('Preflight failed') phase = 'validate' if not stage(directory, 'validate', 1200): raise RuntimeError('Validation failed') if json.loads((directory / 'request.json').read_text())['mode'] == 'plan': status = json.loads((directory / 'status.json').read_text()) status['state'] = 'planned' atomic_json(directory / 'status.json', status) return # Budget includes both rollout checks and rollback waves, plus API overhead. phase = 'Recovery budget' count = int( command( 'bash', str(directory / 'source/.gitea/workflows/deploy-stage.sh'), 'workload-count', env=environment(directory), ) ) verify_budget = max(600, 2 * math.ceil(count / 4) * 300 + 120) if verify_budget > 7200: raise ValueError('More than two hours of recovery required; split this deploy') phase = 'apply-k8s' k8s_ok = stage(directory, 'apply-k8s', 2700) phase = 'apply-compose' compose_ok = stage(directory, 'apply-compose', 1800) if k8s_ok else False phase = 'verify-k8s' verify_ok = stage(directory, 'verify-k8s', verify_budget) phase = 'smoke' smoke_ok = stage(directory, 'smoke', 600) if not all((k8s_ok, compose_ok, verify_ok, smoke_ok)): raise RuntimeError('Deploy failed; inspect stage logs and recovery report') phase = 'Save the successful baseline' finish_success(directory, plan) except Exception as error: status = json.loads((directory / 'status.json').read_text()) status['failure_stage'] = next( (name for name, result in status['stages'].items() if result.get('result') == 'failure'), phase ) atomic_json(directory / 'status.json', status) with (directory / 'controller.log').open('a') as stream: stream.write(f'{error}\n') recover(directory) raise def retain_completed(current): finished = [] for directory in (STATE / 'runs').iterdir(): status_file = directory / 'status.json' if status_file.exists() and json.loads(status_file.read_text())['state'] in ('success', 'planned'): finished.append(directory) for directory in sorted(finished, key=lambda p: p.stat().st_mtime, reverse=True)[20:]: if directory == current: continue command('git', '-C', str(CONFIG_REPO), 'worktree', 'remove', '--force', str(directory / 'source')) shutil.rmtree(directory) def follow(run_id, phase): directory = run_directory(run_id) groups = { 'apply': ('doctor', 'validate', 'apply-k8s', 'apply-compose'), 'verify': ('verify-k8s',), 'smoke': ('smoke',), } names = groups[phase] offsets = {} while True: status = json.loads((directory / 'status.json').read_text()) for name in (*names, 'controller'): path = directory / f'{name}.log' if path.exists(): with path.open() as stream: stream.seek(offsets.get(name, 0)) content = stream.read() if content: print(content, end='', flush=True) offsets[name] = stream.tell() stages = status['stages'] if all(stages.get(name, {}).get('result') in ('success', 'failure') for name in names): return all(stages[name]['result'] == 'success' for name in names) if status['state'] in ('success', 'failure', 'planned'): return status['state'] in ('success', 'planned') time.sleep(3) def summary(run_id): directory = run_directory(run_id) request = json.loads((directory / 'request.json').read_text()) release = request['release'] plan_file = directory / 'plan.json' lines = [ f'## Deploy `{release["sha"]}`', '', f'- Mode: `{request["mode"]}`', f'- Refresh third-party images: `{request["refresh_images"]}`', ] status = json.loads((directory / 'status.json').read_text()) if status.get('failure_stage'): lines.append(f'- Failed stage: **{status["failure_stage"]}**') lines.extend( [ '', f'- Observed run state: **{status["state"]}**', '', '### Stage results', '| Stage | Result | Exit code |', '| --- | --- | --- |', ] ) for name in ('doctor', 'validate', 'apply-k8s', 'apply-compose', 'verify-k8s', 'smoke'): stage_result = status['stages'].get(name, {}) lines.append(f'| {name} | {stage_result.get("result", "not started")} | {stage_result.get("exit_code", "—")} |') lines.extend(['', '### Apply and Helm recovery results']) events_file = directory / 'apply-events.jsonl' events = [] if events_file.exists(): for line in events_file.read_text().splitlines(): try: events.append(json.loads(line)) except json.JSONDecodeError: lines.append('- An operation record is incomplete. Check the stage log.') latest = {(event['action'], event['target']): event['result'] for event in events} lines.extend(f'- `{action}` `{target}`: **{result}**' for (action, target), result in latest.items()) if not latest: lines.append('- No apply results were recorded.') lines.append('- A completed apply does not confirm health. See verification and smoke results.') lines.extend(['', '### Kubernetes recovery']) pointer = directory / 'snapshot/current' failed = Path(pointer.read_text().strip()) / 'failed-workloads' if pointer.exists() else None if failed and failed.exists(): contents = failed.read_text() counts = dict(re.findall(r'^(ROLLED_BACK|UNRECOVERED)=([0-9]+)$', contents, re.MULTILINE)) if not contents.strip(): lines.append('- No failed workloads were recorded. See the verification result above.') elif counts: lines.append(f'- Workloads restored: **{counts.get("ROLLED_BACK", "unknown")}**') lines.append(f'- Workloads that need manual recovery: **{counts.get("UNRECOVERED", "unknown")}**') else: lines.append('- Rollback has no recorded result yet. Check the verification log.') else: lines.append('- No workload rollback was recorded. This does not confirm health.') lines.append('- Compose requires manual recovery. Use the saved command in the apply log.') if not plan_file.exists(): lines.extend(['', 'Plan was not created. Check the controller log.']) print('\n'.join(lines)) return plan = json.loads(plan_file.read_text()) lines.extend(['', '### Selected services']) count = 0 for kind, services in plan['selected'].items(): for service in services: lines.append(f'- `{kind}`: `{service}`') count += 1 if not count: lines.append('- None') lines.extend(['', '### Selected Helm releases']) lines.extend(f'- `{release}`' for release in plan.get('helm', [])) if not plan.get('helm'): lines.append('- None') lines.extend(['', '### Images pinned in the checked release']) lines.extend(f'- `{image}@{digest}`' for image, digest in sorted(release['images'].items())) lines.extend(['', '### Removed resources requiring manual review']) lines.extend(f'- `{item}`' for item in plan.get('removed', [])) if not plan.get('removed'): lines.append('- None') print('\n'.join(lines)) def main(): os.umask(0o077) parser = argparse.ArgumentParser(description=__doc__) parser.add_argument('action', choices=('start', 'execute', 'recover', 'status', 'follow', 'summary')) parser.add_argument('run_id') parser.add_argument('phase', nargs='?', choices=('apply', 'verify', 'smoke')) parser.add_argument('--retry', action='store_true', help='Retry failed recovery checks; never repeat apply') args = parser.parse_args() directory = run_directory(args.run_id) if args.action == 'start': start(args.run_id) elif args.action == 'execute': execute(args.run_id) elif args.action == 'recover': with lock('deploy.lock'): recover(directory, retry=args.retry) elif args.action == 'status': print((directory / 'status.json').read_text()) if (directory / 'plan.json').exists(): plan = json.loads((directory / 'plan.json').read_text()) print(json.dumps({k: plan[k] for k in ('sha', 'selected', 'helm', 'removed')}, indent=2)) elif args.action == 'summary': summary(args.run_id) elif not follow(args.run_id, args.phase): sys.exit(1) if __name__ == '__main__': main()