Files

341 lines
20 KiB
Python

#!/usr/bin/env python3
"""Renew real isolation proofs on an idle PVE deployment; never extend old evidence."""
import argparse
import datetime as dt
import fcntl
import hashlib
import json
import os
from pathlib import Path
import select
import signal
import subprocess
import sys
import time
import uuid
def run(argv, timeout=120, **kwargs):
p = subprocess.run(list(map(str, argv)), text=True, capture_output=True, timeout=timeout, **kwargs)
if p.returncode:
# Arguments/output may contain private operator data. Keep those out of journals.
raise RuntimeError(f'{argv[0]} failed (exit {p.returncode})')
return p.stdout.strip()
def load(path):
return json.loads(Path(path).read_text())
def save(path, value):
path = Path(path)
temporary = path.with_name(path.name + '.new')
with temporary.open('w') as stream:
os.chmod(temporary, 0o600)
json.dump(value, stream, indent=1)
stream.flush()
os.fsync(stream.fileno())
os.replace(temporary, path)
def require(ok, reason):
if not ok:
raise RuntimeError(reason)
class Renewal:
def __init__(self, config):
self.c = config
self.state = Path(config['state_dir'])
self.state.mkdir(mode=0o700, parents=True, exist_ok=True)
self.receipt = self.state / 'paused-worker.json'
self.current = self.state / 'incomplete.json'
self.worker = None
self.work = None
self.clone = None
self.tagged = False
self.binding = None
def ct(self, *args, **kwargs):
return run(['pct', 'exec', str(self.c['control_vmid']), '--', *args], **kwargs)
def docker(self, *args, **kwargs):
return self.ct('docker', *args, **kwargs)
def sql(self, query):
return self.docker('exec', self.c['postgres_container'], 'psql', '-X', '-v', 'ON_ERROR_STOP=1',
'-U', self.c['database_user'], '-d', self.c['database'], '-At', '-c', query)
def recover_pause(self):
if not self.receipt.exists():
return
receipt = load(self.receipt)
info = json.loads(self.docker('inspect', receipt['id']))[0]
require(info['Id'] == receipt['id'] and info['Config']['Labels'].get('com.docker.compose.project') == self.c['project']
and info['Config']['Labels'].get('com.docker.compose.service') == 'worker', 'worker identity changed; manual recovery required')
if info['State']['Paused']:
self.docker('unpause', receipt['id'])
self.receipt.unlink()
print('Original worker unfrozen', flush=True)
def due(self):
config = json.loads(self.ct('cat', self.c['worker_config']))
self.binding = config['owners'][self.c['owner_id']]
self.runtime_config = config
success_path = self.state / 'last-success.json'
if not success_path.exists():
return True
success = load(success_path)
for key, field in ((self.binding['isolation_proof_file'], 'isolation_expires_at'),
(self.binding['online']['proof_file'], 'online_expires_at')):
try:
proof = json.loads(self.ct('cat', self.host_path(key)))
expiry = dt.datetime.fromisoformat(proof['expires_at'].replace('Z', '+00:00'))
published = dt.datetime.fromisoformat(success[field].replace('Z', '+00:00'))
if (proof.get('passed') is not True or proof.get('owner_id') != self.c['owner_id']
or proof.get('node') != self.binding['node'] or abs((expiry - published).total_seconds()) >= .001):
return True
if expiry <= dt.datetime.now(dt.timezone.utc) + dt.timedelta(hours=8):
return True
except (RuntimeError, ValueError, KeyError):
return True
return False
def host_path(self, path):
relative = Path(path).relative_to('/run/otche')
require('..' not in relative.parts, 'unsafe worker mount path')
return str(Path(self.c['worker_config']).parent / relative)
def freeze_idle_worker(self):
# Hold SHARE locks across the idle observation AND freeze. Claims/qualification
# inserts cannot commit between the observation and docker pause.
argv = ['pct', 'exec', str(self.c['control_vmid']), '--', 'docker', 'exec', '-i',
self.c['postgres_container'], 'psql', '-X', '-qAt', '-v', 'ON_ERROR_STOP=1',
'-U', self.c['database_user'], '-d', self.c['database']]
gate = subprocess.Popen(argv, stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.DEVNULL, text=True)
try:
gate.stdin.write("BEGIN; SET LOCAL lock_timeout='2s'; SET LOCAL idle_in_transaction_session_timeout='30s'; LOCK TABLE attempts,qualifications,allocations,extractor_allocations,media IN SHARE MODE NOWAIT;\n"
"SELECT 'idle:' || ((SELECT count(*) FROM attempts WHERE phase<>'finished') + (SELECT count(*) FROM qualifications WHERE status IN ('queued','running')) + (SELECT count(*) FROM allocations WHERE state<>'deleted') + (SELECT count(*) FROM extractor_allocations WHERE state<>'deleted') + (SELECT count(*) FROM media WHERE state<>'deleted'));\n")
gate.stdin.flush()
ready, _, _ = select.select([gate.stdout], [], [], 8)
if not ready or gate.stdout.readline().strip() != 'idle:0':
print('Deferred: work, evidence resources or concurrent database activity', flush=True)
return False
locked_at = time.monotonic()
infos = json.loads(self.docker('inspect', self.c['worker_container'], timeout=5))
info = infos[0]
require(info['State']['Running'] and not info['State']['Paused'], 'worker unavailable or already paused')
require(info['Config']['Labels'].get('com.docker.compose.project') == self.c['project'] and
info['Config']['Labels'].get('com.docker.compose.service') == 'worker', 'wrong worker container')
require(not load('/var/lib/otche-network/state.json')['records'], 'broker has retained leases')
# Durable intent permits ExecStopPost to undo a pause even if killed immediately after it.
save(self.receipt, {'id': info['Id'], 'at': dt.datetime.now(dt.timezone.utc).isoformat()})
self.docker('pause', info['Id'], timeout=8)
self.worker = info['Id']
require(gate.poll() is None and time.monotonic() - locked_at < 25, 'idle gate expired before freeze')
gate.stdin.write("SELECT 'frozen';\n")
gate.stdin.flush()
ready, _, _ = select.select([gate.stdout], [], [], 5)
require(ready and gate.stdout.readline().strip() == 'frozen', 'database freeze transaction lost')
return True
finally:
try:
gate.communicate('ROLLBACK;\n\\q\n', timeout=5)
except subprocess.TimeoutExpired:
gate.kill()
gate.communicate()
def get_vm(self, vmid):
return json.loads(run(['pvesh', 'get', f"/nodes/{self.binding['node']}/qemu/{vmid}/config", '--output-format', 'json']))
def source_projection(self):
source = self.runtime_config['sources'][self.c['source_ref']]
vmid = source['vmid']
require(run(['qm', 'status', str(vmid)]) == 'status: stopped', 'source not stopped')
cfg = self.get_vm(vmid)
require(not cfg.get('lock') and not any(k.startswith('net') for k in cfg), 'source locked or networked')
return cfg
def stage_config(self):
self.work = self.state / (dt.datetime.now(dt.timezone.utc).strftime('%Y%m%dT%H%M%SZ') + '-' + uuid.uuid4().hex[:8])
self.work.mkdir(mode=0o700)
save(self.current, {'work_dir': str(self.work), 'phase': 'staging'})
credentials = self.work / 'credentials'
credentials.mkdir(mode=0o700)
for purpose in ('provisioner', 'runtime', 'recorder', 'uploader', 'housekeeping'):
data = json.loads(self.ct('cat', self.host_path(self.binding[purpose + '_file'])))
save(credentials / (purpose + '.json'), data)
network = dict(self.binding['online'])
for key in ('ca_file', 'certificate_file', 'key_file'):
target = credentials / ('network-' + key)
target.write_text(self.ct('cat', self.host_path(network[key])) + '\n')
target.chmod(0o600)
network[key] = str(target)
save(self.work / 'worker.json', self.runtime_config)
source = self.runtime_config['sources'][self.c['source_ref']]
self.probe = {'owner_id': self.c['owner_id'], 'attempt_id': str(uuid.uuid4()), 'node': self.binding['node'],
'source_vmid': source['vmid'], 'backend_dir': self.c['backend_dir'], 'worker_config': str(self.work / 'worker.json'),
'credentials_dir': str(credentials), 'pve_ca_file': '/etc/pve/pve-root-ca.pem', 'network': network,
'broker_config_file': self.c['broker_config_file'], 'broker_service': self.c['broker_service'],
'network_probe': self.c['network_probe']}
save(self.work / 'probe-config.json', self.probe)
def create_clone(self):
before = self.source_projection()
node = json.loads(run(['pvesh', 'get', '/nodes/' + self.binding['node'] + '/status', '--output-format', 'json']))
require(node['memory']['free'] >= self.binding['min_memory_free_bytes'] + (int(before['memory']) + 1024) * 1024**2, 'probe RAM headroom insufficient')
storage = json.loads(run(['pvesh', 'get', '/nodes/' + self.binding['node'] + '/storage/' + self.binding['disk_storage'] + '/status', '--output-format', 'json']))
source = self.runtime_config['sources'][self.c['source_ref']]
require(storage['avail'] >= self.binding['min_storage_free_bytes'] + source['max_disk_bytes'], 'probe disk headroom insufficient')
vmid = int(json.loads(run(['pvesh', 'get', '/cluster/nextid', '--output-format', 'json'])))
require(vmid not in {7000, 7001, self.c['control_vmid']} | {v['vmid'] for v in self.runtime_config['sources'].values()}, 'protected VMID')
name = 'otche-' + self.probe['attempt_id']
save(self.work / 'clone.json', {'vmid': vmid, 'name': name, 'phase': 'intent', 'source_config': before})
run(['qm', 'clone', str(source['vmid']), str(vmid), '--full', '1', '--pool', self.binding['pool'],
'--storage', self.binding['disk_storage'], '--name', name], timeout=600)
self.clone = vmid
cfg = self.get_vm(vmid)
require(cfg['name'] == name and not cfg.get('lock') and not any(k.startswith('net') for k in cfg), 'clone identity mismatch')
run(['qm', 'set', str(vmid), '--digest', cfg['digest'], '--tags', 'otche;otche-owner-' + self.c['owner_id'] + ';otche-attempt-' + self.probe['attempt_id'], '--onboot', '0'])
self.tagged = True
save(self.work / 'clone.json', {'vmid': vmid, 'name': name, 'phase': 'tagged', 'source_config': before})
require(self.source_projection() == before, 'source changed during clone')
self.source_before = before
def delete_clone(self):
if self.clone is None:
return
require(self.tagged, 'ambiguous untagged clone; inspect retained clone receipt')
cfg = self.get_vm(self.clone)
tags = set(cfg.get('tags', '').split(';'))
require(cfg['name'] == 'otche-' + self.probe['attempt_id'] and {'otche', 'otche-owner-' + self.c['owner_id'],
'otche-attempt-' + self.probe['attempt_id']} <= tags and not cfg.get('lock'), 'refusing changed probe deletion')
require(run(['qm', 'status', str(self.clone)]) == 'status: stopped' and not any(k.startswith('net') for k in cfg), 'probe must be stopped and NIC-free')
pool = json.loads(run(['pvesh', 'get', '/pools/' + self.binding['pool'], '--output-format', 'json']))
require(any(m.get('vmid') == self.clone and m.get('node') == self.binding['node'] for m in pool['members']), 'probe pool mismatch')
run(['qm', 'destroy', str(self.clone), '--purge', '0', '--destroy-unreferenced-disks', '0'], timeout=120)
save(self.work / 'clone-cleanup.json', {'vmid': self.clone, 'deleted': True})
self.clone = None
def publish(self):
pairs = [('owner', self.binding['isolation_proof_file']), ('online', self.binding['online']['proof_file'])]
# Validate both bundles and freshness before modifying either live proof.
for name, destination in pairs:
proof = load(self.work / (name + '-proof.json'))
evidence = (self.work / (name + '-evidence.json')).read_bytes()
require(proof['passed'] and proof['owner_id'] == self.c['owner_id'] and proof['node'] == self.binding['node']
and proof['evidence_sha256'] == hashlib.sha256(evidence).hexdigest(), 'invalid new proof bundle')
expiry = dt.datetime.fromisoformat(proof['expires_at'].replace('Z', '+00:00'))
now = dt.datetime.now(dt.timezone.utc)
require(now + dt.timedelta(hours=23) < expiry <= now + dt.timedelta(hours=24), 'proof not newly measured')
for name, destination in pairs:
dest = self.host_path(destination)
# Preserve immutable evidence by content-addressed name, atomically publish proof last.
digest = load(self.work / (name + '-proof.json'))['evidence_sha256']
evidence_dest = str(Path(dest).parent / (name + '-evidence-' + digest + '.json'))
for local, target in [(self.work / (name + '-evidence.json'), evidence_dest), (self.work / (name + '-proof.json'), dest)]:
staged = target + '.renew-' + self.probe['attempt_id']
run(['pct', 'push', str(self.c['control_vmid']), str(local), staged, '--user', '10001', '--group', '10001', '--perms', '0600'])
self.ct('mv', staged, target)
# Exec inside the existing approved image; no builds, pulls, migrations or project changes.
info = json.loads(self.docker('inspect', self.worker))[0]
image = info['Image']
command = ['run', '--rm', '--network', 'container:' + self.worker, '--volumes-from', self.worker + ':ro',
'--env', 'DATABASE_URL_FILE=/run/secrets/database-url', '--env', 'WORKER_CONFIG_FILE=/run/otche/config.json',
image, 'bindings-sync']
self.docker(*command, timeout=120)
row = json.loads(self.sql("SELECT row_to_json(b) FROM bindings b WHERE owner_id='" + str(uuid.UUID(self.c['owner_id'])) + "'"))
require(row['configured'] and row['online_ready'], 'bindings-sync did not publish readiness')
for name, field in [('owner', 'isolation_expires_at'), ('online', 'online_expires_at')]:
actual = dt.datetime.fromisoformat(row[field].replace('Z', '+00:00'))
wanted = dt.datetime.fromisoformat(load(self.work / (name + '-proof.json'))['expires_at'].replace('Z', '+00:00'))
require(abs((actual - wanted).total_seconds()) < .001, 'binding expiry differs from fresh proof')
save(self.state / 'last-success.json', {'at': dt.datetime.now(dt.timezone.utc).isoformat(), 'work_dir': str(self.work),
'isolation_expires_at': row['isolation_expires_at'], 'online_expires_at': row['online_expires_at']})
def execute(self, force=False):
require(not self.current.exists(), 'previous incomplete probe requires operator review; see incomplete.json')
if not self.due() and not force:
print('Proofs remain valid beyond the renewal window', flush=True)
return
for directory, commit in [(self.c['backend_dir'], self.c['backend_commit']),
(str(Path(__file__).resolve().parent), self.c['deploy_commit'])]:
require(run(['git', '-C', directory, 'rev-parse', 'HEAD']) == commit and
not run(['git', '-C', directory, 'status', '--porcelain', '--untracked-files=no']),
'automation source differs from approved clean commit')
try:
if not self.freeze_idle_worker():
return
self.stage_config()
self.create_clone()
for name in ('owner', 'network'):
with (self.work / (name + '.log')).open('w') as log:
p = subprocess.run(['python3', str(Path(self.c['backend_dir']) / 'provisioning' / ('revalidate-' + name + '.py')),
'--config', str(self.work / 'probe-config.json'), '--work-dir', str(self.work), '--vmid', str(self.clone)],
stdout=log, stderr=subprocess.STDOUT, timeout=1500)
require(p.returncode == 0, name + ' probes failed; see private run log')
require(self.source_projection() == self.source_before, 'source changed during revalidation')
self.delete_clone()
require(not load('/var/lib/otche-network/state.json')['records'], 'broker records remain after probes')
self.publish()
self.current.unlink()
print('Fresh owner and online proofs published after real probes and cleanup', flush=True)
except Exception:
if self.work:
save(self.state / 'last-failure.json', {'at': dt.datetime.now(dt.timezone.utc).isoformat(), 'work_dir': str(self.work)})
owner_journal = self.work / 'owner-resources.json'
network_journal = self.work / 'network-resources.json'
owner_clean = owner_journal.exists() and any(row['phase'] == 'complete' and row['resource'] == 'cleanup'
for row in load(owner_journal)['events'])
network_clean = not network_journal.exists() or load(network_journal).get('cleanup_complete') is True
if not (self.work / 'clone.json').exists() and not owner_journal.exists() and not network_journal.exists():
# No external mutation intent exists: capacity/config reads can be retried.
self.current.unlink(missing_ok=True)
elif owner_clean and network_clean and not load('/var/lib/otche-network/state.json')['records']:
self.delete_clone()
self.current.unlink(missing_ok=True)
print('Failed checks left no resources; next timer may retry without renewing old proofs', flush=True)
raise
finally:
# A failed/ambiguous test retains its journal and cannot be automatically overwritten.
# Never force-delete a running/changed VM or hide a failed cleanup.
self.recover_pause()
if self.work:
for path in (self.work / 'credentials').glob('*'):
path.unlink()
def main():
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument('--config', required=True)
parser.add_argument('--force', action='store_true', help='Run even outside renewal window; never bypass idle gate')
parser.add_argument('--recover-worker', action='store_true', help='ExecStopPost: thaw only the recorded exact worker')
args = parser.parse_args()
require(os.geteuid() == 0, 'Linux root required')
os.umask(0o077)
cfgpath = Path(args.config)
require(cfgpath.stat().st_uid == 0 and cfgpath.stat().st_mode & 0o077 == 0, 'private root config required')
renewal = Renewal(load(cfgpath))
with (renewal.state / 'renew.lock').open('a') as lock:
try:
fcntl.flock(lock, fcntl.LOCK_EX | fcntl.LOCK_NB)
except BlockingIOError:
print('Another revalidation owns the lock', flush=True)
return
if args.recover_worker:
renewal.recover_pause()
return
def interrupted(signum, frame):
raise RuntimeError('renewal interrupted')
signal.signal(signal.SIGTERM, interrupted)
signal.signal(signal.SIGINT, interrupted)
renewal.execute(args.force)
if __name__ == '__main__':
try:
main()
except Exception as error:
print(str(error), file=sys.stderr, flush=True)
sys.exit(1)