341 lines
20 KiB
Python
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)
|