From 805e0e5a2b6263e4601c39230eb76d2f1bc7243a Mon Sep 17 00:00:00 2001 From: codex Date: Sun, 6 Sep 2026 23:46:20 +0200 Subject: [PATCH] Verify telemetry essentials recovery from versioned Scaleway archive Assistant: codex Assistant-Model: gpt-6-astra Assistant-Session: 01a06ecb-456a-71c2-b41e-0755d336e883 --- .../telemetry-grafana.md | 5 +- .../telemetry-capture-2026-09-06.json | 18 ++ .../telemetry-primary-2026-09-06.json | 12 + ...telemetry-recovery-cleanup-2026-09-06.json | 5 + .../telemetry-restore-2026-09-06.json | 61 ++++ .../2026-09-06-telemetry-scaleway-recovery.md | 63 ++++ scripts/restore_telemetry_essentials.py | 274 ++++++++++++++++++ scripts/telemetry_primary_archive.py | 85 ++++++ tests/test_telemetry_primary_archive.py | 40 +++ 9 files changed, 561 insertions(+), 2 deletions(-) create mode 100644 docs/evidence/telemetry-capture-2026-09-06.json create mode 100644 docs/evidence/telemetry-primary-2026-09-06.json create mode 100644 docs/evidence/telemetry-recovery-cleanup-2026-09-06.json create mode 100644 docs/evidence/telemetry-restore-2026-09-06.json create mode 100644 history/2026-09-06-telemetry-scaleway-recovery.md create mode 100644 scripts/restore_telemetry_essentials.py create mode 100644 scripts/telemetry_primary_archive.py create mode 100644 tests/test_telemetry_primary_archive.py diff --git a/docs/credential-lane-designs/telemetry-grafana.md b/docs/credential-lane-designs/telemetry-grafana.md index 9ba320a..9989562 100644 --- a/docs/credential-lane-designs/telemetry-grafana.md +++ b/docs/credential-lane-designs/telemetry-grafana.md @@ -40,5 +40,6 @@ Do not delete retained recovery material as an implicit rollback. Accepted September 6 evidence: exact read and delivery match; wrong identity, namespace and audience denied; sibling/write denied; coding-agent deny wins; native login succeeds; anonymous/proxy-header requests denied. See -../rapp-telemetry/evidence/live/2026-09-06-railiance01.json. Independent backup and -isolated restore, operator OIDC and public production admission remain open. +../rapp-telemetry/evidence/live/2026-09-06-railiance01.json. A version-pinned Scaleway essentials archive and isolated application restore +passed later on September 6; see history/2026-09-06-telemetry-scaleway-recovery.md. +Recurring backups, operator OIDC and public production admission remain open. diff --git a/docs/evidence/telemetry-capture-2026-09-06.json b/docs/evidence/telemetry-capture-2026-09-06.json new file mode 100644 index 0000000..07b4d8e --- /dev/null +++ b/docs/evidence/telemetry-capture-2026-09-06.json @@ -0,0 +1,18 @@ +{ + "schema": "rapp-telemetry.capture.v1", + "status": "captured", + "stage": "encrypted", + "helper": "telemetry-capture-018fd2cb75", + "images": { + "grafana": "docker.io/grafana/grafana@sha256:3600073f45a9895d2e2c4042da0e745819f5b78bf781266952a400c40dcd8d69", + "alertmanager": "quay.io/prometheus/alertmanager@sha256:690c7b525f4367aa91f73e2f91c632206d32e97c6384bdbf2fb7a861b420340d" + }, + "source_commit": "37b663953c9522665481ce75fbb5e37f066e55e6", + "fence_started_at": "2026-09-06T20:53:38.480651+00:00", + "writers_stopped": true, + "fence_ended_at": "2026-09-06T20:53:46.007312+00:00", + "ciphertext_bytes": 3748952, + "ciphertext_sha256": "c039d618674b6fe4bd6815d4fa1f3dbf2f52e845c4af3bd52cffdd39e14d4134", + "cleanup_errors": [], + "production_resumed": true +} diff --git a/docs/evidence/telemetry-primary-2026-09-06.json b/docs/evidence/telemetry-primary-2026-09-06.json new file mode 100644 index 0000000..b78b9e1 --- /dev/null +++ b/docs/evidence/telemetry-primary-2026-09-06.json @@ -0,0 +1,12 @@ +{ + "schema": "platform.telemetry-primary-archive.v1", + "status": "fetched_pending_restore", + "started_at": "2026-09-06T20:54:42.876593+00:00", + "destination": "s3://railiance-platform-pg-backup/platform-pg/application-archives/telemetry/2026/09/06/205444-17fc36f5db9e0c198b6681290f02fdc3.tar.age", + "ciphertext_sha256": "c039d618674b6fe4bd6815d4fa1f3dbf2f52e845c4af3bd52cffdd39e14d4134", + "ciphertext_bytes": 3748952, + "version_id": "1788728086280183", + "version_pinned": true, + "download_hash_matches": true, + "finished_at": "2026-09-06T20:54:45.698430+00:00" +} \ No newline at end of file diff --git a/docs/evidence/telemetry-recovery-cleanup-2026-09-06.json b/docs/evidence/telemetry-recovery-cleanup-2026-09-06.json new file mode 100644 index 0000000..d6ccbb9 --- /dev/null +++ b/docs/evidence/telemetry-recovery-cleanup-2026-09-06.json @@ -0,0 +1,5 @@ +{ + "fixture_expired": true, + "targets": 12, + "targets_up": 12 +} \ No newline at end of file diff --git a/docs/evidence/telemetry-restore-2026-09-06.json b/docs/evidence/telemetry-restore-2026-09-06.json new file mode 100644 index 0000000..d9e5d21 --- /dev/null +++ b/docs/evidence/telemetry-restore-2026-09-06.json @@ -0,0 +1,61 @@ +{ + "schema": "platform.telemetry-recovery.v1", + "status": "verified", + "started_at": "2026-09-06T21:38:41.340137+00:00", + "stage": "application_acceptance", + "source_sha256": "c039d618674b6fe4bd6815d4fa1f3dbf2f52e845c4af3bd52cffdd39e14d4134", + "source_version": "1788728086280183", + "images": { + "grafana": "docker.io/grafana/grafana@sha256:3600073f45a9895d2e2c4042da0e745819f5b78bf781266952a400c40dcd8d69", + "alertmanager": "quay.io/prometheus/alertmanager@sha256:690c7b525f4367aa91f73e2f91c632206d32e97c6384bdbf2fb7a861b420340d" + }, + "decrypted": true, + "sqlite_integrity": "ok", + "legacy_database_counts": { + "dashboards": 0, + "users": 1, + "datasources": 2 + }, + "provisioning_files_restored": 24, + "isolation": "docker internal network; no published ports; internal probe; null notification receiver", + "applications_ready": true, + "anonymous_denied": true, + "proxy_headers_denied": true, + "restored_dashboard_count": 20, + "restored_datasource_count": 2, + "native_admin_login": true, + "dashboard_and_datasource_counts_match": true, + "alertmanager_fixture_restored": true, + "source_commit": "37b663953c9522665481ce75fbb5e37f066e55e6", + "startup_diagnostics": [ + { + "component": "probe", + "running": true, + "exit_code": 0, + "oom_killed": false, + "flags": [] + }, + { + "component": "grafana", + "running": true, + "exit_code": 0, + "oom_killed": false, + "flags": [ + "provisioning", + "migration", + "no such file", + "error=" + ] + }, + { + "component": "alertmanager", + "running": true, + "exit_code": 0, + "oom_killed": false, + "flags": [] + } + ], + "cleanup_errors": [], + "plaintext_removed": true, + "finished_at": "2026-09-06T21:39:02.475104+00:00" +} diff --git a/history/2026-09-06-telemetry-scaleway-recovery.md b/history/2026-09-06-telemetry-scaleway-recovery.md new file mode 100644 index 0000000..d608af3 --- /dev/null +++ b/history/2026-09-06-telemetry-scaleway-recovery.md @@ -0,0 +1,63 @@ +# Telemetry Scaleway recovery — 2026-09-06 + +The user authorized the next telemetry backup/recovery step. The private service +now has a verified essentials archive on Scaleway and an isolated application +restore. RAPP-TELEMETRY-WP-0001-T03 can close for private installation/recovery +proof. This does not grant public production admission or close RPF-WP-0036. + +## Accepted evidence + +- 3,748,952 bytes encrypted with the established Railiance age recipient. +- Prefix platform-pg/application-archives/telemetry/ in the existing Scaleway + railiance-platform-pg-backup bucket, separate from Barman/Forgejo. Exact version + 1788728086280183 was downloaded and SHA-256 matched the source ciphertext. +- Preserves full Grafana data/plugins, Alertmanager persistent snapshots, + ConfigMaps, Alertmanager configuration and a source Git archive. Omits disposable + Prometheus TSDB; production rebuilds its seven-day metrics history. No telemetry + upload to Nextcloud occurred. +- Grafana and Alertmanager replicas were fenced for 7.53 seconds for a consistent + volume copy, then both rollout checks passed. This is the replica-fence window, + not a measured endpoint-availability SLA. Final production health is 12/12 + scrape targets up. +- The fetched copy was decrypted with existing governed escrow in a contained + attended session. Same-digest Grafana and Alertmanager started on an internal + local Docker network, no published ports, with alert delivery disabled. +- SQLite integrity passed; native admin login worked; anonymous and forged-header + requests were denied; 20 dashboards, two datasources and the exact inert + Alertmanager silence were recovered. +- Temporary containers, networks and plaintext were removed. The inert production + silence was expired and verified. The encrypted Scaleway recovery point remains. + +## Problems resolved during the drill + +The first content probe assumed Grafana's legacy dashboard table was authoritative. +Grafana 13 migrated dashboards to unified storage; API acceptance is the correct +check, and the legacy table count of zero was not evidence of loss. Next, Docker's +internal network did not provide the assumed host port bindings, so the probe now +runs inside that network with no published ports. Grafana required its production +emptyDir search mount to be reconstructed as writable temporary storage. Finally, +its archived provisioning ConfigMaps needed to be mounted and its search index +allowed to rebuild. The final unmodified run of the corrected tool passed in +about 21 seconds with images cached; do not call this a full disaster-recovery RTO. +All failed attempts were isolated and cleaned up; no production recovery rollback +was performed. Fixed diagnostic classifications contain no secret log output. + +The credential catalog id backup-object-storage is unresolved, as the existing +CCR-2026-0012 already records. Transfer used that CCR's established databases +Secret delivery within the platform owner tool; no new namespace grant or +workload key copy was created. Decryption used the established attended operator +route, and each session self-revoked. No credential values entered receipts. + +## Remaining production work + +Bind daily capture/upload and missed-run reporting to a durable, scoped executor; +these tools establish an attended repeatable path, not an automatic backup SLA. +Keep independent recovery-key/unseal availability explicit: reading current +OpenBao during this drill is not a total-cluster-loss escrow proof. Existing +provider retention remains owned by reef-storage/resource-control; no lifecycle +or deletion policy was changed. Nextcloud remains essentials-only, with no TSDB. + +Operator SSO, actual signal delivery/receipt, an outside-node watchdog and public +production admission remain under T04. Package evidence is in +../rapp-telemetry/evidence/live/2026-09-06-isolated-restore.json and procedures in +../rapp-telemetry/docs/recovery.md. diff --git a/scripts/restore_telemetry_essentials.py b/scripts/restore_telemetry_essentials.py new file mode 100644 index 0000000..bfadd59 --- /dev/null +++ b/scripts/restore_telemetry_essentials.py @@ -0,0 +1,274 @@ +#!/usr/bin/env python3 +"""Silent attended recovery from a verified Scaleway download into isolated Docker.""" +import argparse +import base64 +from datetime import datetime, timezone +import hashlib +import importlib.util +import json +import os +from pathlib import Path +import secrets +import sqlite3 +import subprocess +import tarfile +import tempfile +import time +import urllib.error +import urllib.request +from state_hub_preflight_lane import bao, data, LaneError + +PROBE_CONTAINER = None +PROBE_IMAGE = 'python@sha256:78387bc3881b8273120a12ebe6c1ab22b018ccc2c9adf565ae1ac9b536e184ea' + +PACKAGE = Path(__file__).resolve().parents[2] / 'rapp-telemetry' +spec = importlib.util.spec_from_file_location('capture', PACKAGE / 'tools/capture_recovery.py') +capture = importlib.util.module_from_spec(spec) +spec.loader.exec_module(capture) + + +def cmd(argv, **kwargs): + result = subprocess.run(argv, capture_output=True, timeout=90, **kwargs) + if result.returncode: + raise LaneError('command_failed') + return result.stdout + + +class NoRedirect(urllib.request.HTTPRedirectHandler): + def redirect_request(self, *args, **kwargs): + return None + + +PROBE_SCRIPT = """import json,sys,urllib.request,urllib.error +class NoRedirect(urllib.request.HTTPRedirectHandler): + def redirect_request(self,*args,**kwargs):return None +p=json.load(sys.stdin) +r=urllib.request.Request(p['url'],headers=p['headers']) +try: + with urllib.request.build_opener(urllib.request.ProxyHandler({}),NoRedirect()).open(r,timeout=3) as out: + print(json.dumps({'code':out.status,'body':out.read(4194304).decode()})) +except urllib.error.HTTPError as e:print(json.dumps({'code':e.code,'body':''})) +except Exception:print(json.dumps({'error':'probe_unavailable'})) +""" + + +def probe(base, path, headers=None): + if PROBE_CONTAINER is None: + raise LaneError('isolated_probe_required') + result = json.loads(cmd(['docker', 'exec', '-i', PROBE_CONTAINER, 'python3', '-c', PROBE_SCRIPT], input=json.dumps({'url': base + path, 'headers': headers or {}}).encode())) + if 'error' in result: + raise OSError('probe_unavailable') + return result['code'], result['body'].encode() + + +def wait_health(base, path): + for _ in range(90): + try: + if probe(base, path)[0] == 200: + return + except (OSError, urllib.error.URLError): + pass + time.sleep(1) + raise LaneError('application_not_ready') + + +def run(args, receipt, checkpoint): + global PROBE_CONTAINER + transfer = json.loads(args.transfer_receipt.read_text()) + if transfer.get('status') != 'fetched_pending_restore' or not transfer.get('version_pinned') or not transfer.get('download_hash_matches'): + raise LaneError('verified_primary_download_required') + with args.source.open('rb') as source: + digest = hashlib.file_digest(source, 'sha256').hexdigest() + if digest != transfer['ciphertext_sha256']: + raise LaneError('download_changed') + policies = data(bao(['token', 'lookup', '-format=json']))['data']['policies'] + if 'platform-admin' not in policies or 'root' in policies: + raise LaneError('attended_operator_required') + receipt.update(stage='escrow_decryption', source_sha256=digest, source_version=transfer['version_id']) + checkpoint() + key = data(bao(['read', '-format=json', 'platform/data/workloads/railiance/backup/offsite-lane']))['data']['data']['AGE_PRIVATE_KEY'] + key_bytes = (key.strip() + '\n').encode() + recipient = cmd(['age-keygen', '-y'], input=key_bytes).decode().strip() + if recipient != capture.RECIPIENT: + raise LaneError('escrow_recipient_mismatch') + with tempfile.TemporaryDirectory(prefix='telemetry-restore-') as directory: + directory = Path(directory) + plain = directory / 'fetched.tar' + with plain.open('xb') as out: + result = subprocess.run(['age', '-d', '-i', '/dev/stdin', str(args.source)], input=key_bytes, stdout=out, stderr=subprocess.PIPE, timeout=90) + if result.returncode: + raise LaneError('decryption_failed') + root = directory / 'restored' + root.mkdir(mode=0o700) + with tarfile.open(plain) as archive: + capture.validate_members(archive) + archive.extractall(root, filter='data') + manifest = json.loads((root / 'recovery-manifest.json').read_text()) + if manifest['schema'] != 'rapp-telemetry.essentials.v1': + raise LaneError('wrong_manifest') + images = manifest['images'] + for name in ['grafana', 'alertmanager']: + if '@sha256:' not in images[name]: + raise LaneError('unpinned_image') + cmd(['docker', 'image', 'inspect', images[name]]) + receipt.update(stage='database_validation', images=images, decrypted=True) + checkpoint() + with sqlite3.connect('file:' + str(root / 'grafana/grafana.db') + '?mode=ro', uri=True) as db: + if db.execute('PRAGMA integrity_check').fetchone()[0] != 'ok': + raise LaneError('sqlite_integrity_failed') + counts = {'dashboards': db.execute('SELECT count(*) FROM dashboard WHERE is_folder=0').fetchone()[0], 'users': db.execute('SELECT count(*) FROM user').fetchone()[0], 'datasources': db.execute('SELECT count(*) FROM data_source').fetchone()[0]} + receipt.update(sqlite_integrity='ok', legacy_database_counts=counts) + if counts['users'] < 1: + raise LaneError('required_user_database_content_missing') + configs = json.loads((root / 'config/configmaps.json').read_text())['items'] + grafana_config = next(c['data']['grafana.ini'] for c in configs if c['metadata']['name'] == 'telemetry-grafana') + config = directory / 'grafana.ini' + config.write_text(grafana_config) + config.chmod(0o600) + provisioning = directory / 'provisioning' + dashboards = directory / 'dashboards' + dashboards.mkdir() + for part in ['dashboards', 'datasources']: + (provisioning / part).mkdir(parents=True) + for item in configs: + labels = item['metadata'].get('labels', {}) + target = None + if item['metadata']['name'] == 'telemetry-grafana-config-dashboards': + target = provisioning / 'dashboards' + elif labels.get('grafana_dashboard') == '1': + target = dashboards + elif labels.get('grafana_datasource') == '1': + target = provisioning / 'datasources' + if target: + for filename, content in item.get('data', {}).items(): + if Path(filename).name != filename: + raise LaneError('unsafe_config_filename') + (target / filename).write_text(content) + receipt['provisioning_files_restored'] = len(list(provisioning.rglob('*'))) + len(list(dashboards.iterdir())) + # Notification sinks are intentionally absent in the recovery environment. + am_config = directory / 'alertmanager.yaml' + am_config.write_text('route:\n receiver: restore-null\nreceivers:\n - name: restore-null\n') + state = list((root / 'alertmanager').rglob('silences')) + if len(state) != 1 or not (state[0].parent / 'nflog').is_file(): + raise LaneError('alertmanager_snapshots_missing') + network = 'telemetry-recovery-' + secrets.token_hex(6) + containers = [] + network_created = False + receipt.update(stage='isolated_startup', isolation='docker internal network; no published ports; internal probe; null notification receiver') + checkpoint() + try: + cmd(['docker', 'network', 'create', '--internal', network]) + network_created = True + if not json.loads(cmd(['docker', 'network', 'inspect', network]))[0]['Internal']: + raise LaneError('network_not_internal') + def start(name, image, port, mounts, extra, env=()): + container = network + '-' + name + argv = ['docker', 'create', '--name', container, '--network', network, '--user', str(os.getuid()) + ':' + str(os.getgid()), '--read-only', '--tmpfs', '/tmp:rw,nosuid,size=64m', '--tmpfs', '/var/lib/grafana-search:rw,nosuid,size=128m,uid=' + str(os.getuid()) + ',gid=' + str(os.getgid()), '--cap-drop=ALL', '--security-opt=no-new-privileges', '--pids-limit=256', '--memory=512m', '--cpus=0.5'] + for source, target, mode in mounts: + argv += ['--volume', str(source) + ':' + target + ':' + mode] + for item in env: + argv += ['--env', item] + cmd(argv + [image] + extra) + containers.append(container) + cmd(['docker', 'start', container]) + info = json.loads(cmd(['docker', 'inspect', container]))[0] + if info['HostConfig'].get('PortBindings') or set(info['NetworkSettings']['Networks']) != {network}: + raise LaneError('recovery_network_mismatch') + return 'http://' + container + ':' + str(port) + probe_name = network + '-probe' + cmd(['docker', 'create', '--name', probe_name, '--network', network, '--read-only', '--cap-drop=ALL', '--security-opt=no-new-privileges', '--memory=64m', '--pids-limit=64', PROBE_IMAGE, 'python3', '-c', 'import time; time.sleep(600)']) + containers.append(probe_name) + cmd(['docker', 'start', probe_name]) + PROBE_CONTAINER = probe_name + grafana = start('grafana', images['grafana'], 3000, [(root / 'grafana', '/var/lib/grafana', 'rw'), (config, '/etc/grafana/grafana.ini', 'ro'), (provisioning, '/etc/grafana/provisioning', 'ro'), (dashboards, '/tmp/dashboards', 'ro')], [], ['GF_PATHS_DATA=/var/lib/grafana', 'GF_PATHS_LOGS=/tmp/grafana-logs', 'GF_SERVER_HTTP_ADDR=0.0.0.0', 'GF_SERVER_DOMAIN=localhost', 'GF_SERVER_ROOT_URL=http://localhost:3000', 'GF_UNIFIED_ALERTING_ENABLED=false', 'GF_ANALYTICS_REPORTING_ENABLED=false', 'GF_PLUGINS_PREINSTALL_DISABLED=true']) + alertmanager = start('alertmanager', images['alertmanager'], 9093, [(state[0].parent, '/data', 'rw'), (am_config, '/restore.yaml', 'ro')], ['--config.file=/restore.yaml', '--storage.path=/data', '--cluster.listen-address=', '--web.listen-address=0.0.0.0:9093']) + receipt['stage'] = 'grafana_startup' + checkpoint() + wait_health(grafana, '/api/health') + receipt['stage'] = 'alertmanager_startup' + checkpoint() + wait_health(alertmanager, '/-/ready') + receipt.update(stage='application_acceptance', applications_ready=True) + checkpoint() + for label, headers in [('anonymous_denied', {}), ('proxy_headers_denied', {'X-WEBAUTH-USER': 'admin', 'Remote-User': 'admin'})]: + if probe(grafana, '/api/user', headers)[0] != 401: + raise LaneError('negative_login_failed') + receipt[label] = True + native = data(bao(['read', '-format=json', 'platform/data/workloads/telemetry/grafana-admin']))['data']['data'] + auth = base64.b64encode((native['ADMIN_USERNAME'] + ':' + native['ADMIN_PASSWORD']).encode()).decode() + headers = {'Authorization': 'Basic ' + auth} + code, body = probe(grafana, '/api/user', headers) + if code != 200 or not json.loads(body).get('isGrafanaAdmin'): + raise LaneError('restored_admin_login_failed') + for _ in range(90): + code, body = probe(grafana, '/api/search?type=dash-db&limit=1000', headers) + receipt['restored_dashboard_count'] = len(json.loads(body)) if code == 200 else None + if receipt['restored_dashboard_count'] == args.expected_dashboard_count: + break + time.sleep(1) + else: + raise LaneError('restored_dashboard_count_mismatch') + code, body = probe(grafana, '/api/datasources', headers) + if code != 200: + raise LaneError('restored_datasource_read_failed') + sources = json.loads(body) + if not any(s['type'] == 'prometheus' for s in sources) or counts['datasources'] and len(sources) != counts['datasources']: + raise LaneError('restored_datasource_mismatch') + receipt['restored_datasource_count'] = len(sources) + code, body = probe(alertmanager, '/api/v2/silence/' + manifest['silence_id']) + if code != 200 or json.loads(body)['id'] != manifest['silence_id']: + raise LaneError('restored_silence_missing') + receipt.update(status='verified', native_admin_login=True, dashboard_and_datasource_counts_match=True, alertmanager_fixture_restored=True, source_commit=manifest['source_commit']) + finally: + cleanup_errors = [] + receipt['startup_diagnostics'] = [] + for item in containers: + try: + state_info = json.loads(cmd(['docker', 'inspect', item]))[0]['State'] + log_result = subprocess.run(['docker', 'logs', '--tail', '100', item], capture_output=True, timeout=10) + logs = (log_result.stdout + log_result.stderr).decode(errors='replace').lower() + receipt['startup_diagnostics'].append({'component': item.rsplit('-', 1)[-1], 'running': state_info['Running'], 'exit_code': state_info['ExitCode'], 'oom_killed': state_info['OOMKilled'], 'flags': [flag for flag in ['permission denied', 'read-only file system', 'database is locked', 'unable to open', 'error loading config', 'failed to start', 'provisioning', 'migration', 'panic', 'address already in use', 'no such file', 'not found', 'read-only', 'is not writable', 'error=','error:'] if flag in logs]}) + except Exception: + receipt['startup_diagnostics'].append({'component': item.rsplit('-', 1)[-1], 'diagnostic_unavailable': True}) + for container in reversed(containers): + try: + cmd(['docker', 'rm', '-f', container]) + except Exception: + cleanup_errors.append('container_cleanup_failed') + if network_created: + try: + cmd(['docker', 'network', 'rm', network]) + except Exception: + cleanup_errors.append('network_cleanup_failed') + PROBE_CONTAINER = None + receipt['cleanup_errors'] = cleanup_errors + if cleanup_errors: + raise LaneError('recovery_cleanup_failed') + receipt['plaintext_removed'] = True + + +def main(): + p = argparse.ArgumentParser(description=__doc__) + for name in ['source', 'transfer-receipt', 'receipt']: + p.add_argument('--' + name, type=Path, required=True) + p.add_argument('--expected-dashboard-count', type=int, required=True) + args = p.parse_args() + if args.expected_dashboard_count < 1: + p.error('positive baseline dashboard count required') + fd = os.open(args.receipt, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600) + os.close(fd) + receipt = {'schema': 'platform.telemetry-recovery.v1', 'status': 'failed', 'started_at': datetime.now(timezone.utc).isoformat()} + def checkpoint(): + args.receipt.write_text(json.dumps(receipt, indent=2) + '\n') + try: + run(args, receipt, checkpoint) + except Exception as error: + receipt.update(status='failed', error_type=type(error).__name__, error=str(error) if isinstance(error, LaneError) else 'recovery_failed') + finally: + receipt['finished_at'] = datetime.now(timezone.utc).isoformat() + checkpoint() + return int(receipt['status'] != 'verified') + + +if __name__ == '__main__': + raise SystemExit(main()) diff --git a/scripts/telemetry_primary_archive.py b/scripts/telemetry_primary_archive.py new file mode 100644 index 0000000..437341f --- /dev/null +++ b/scripts/telemetry_primary_archive.py @@ -0,0 +1,85 @@ +#!/usr/bin/env python3 +"""Create and re-fetch a bounded encrypted telemetry object; no delete/list operations.""" +import argparse +import base64 +from datetime import datetime, timezone +import hashlib +import json +import os +from pathlib import Path +import secrets +from state_hub_preflight_lane import assert_cluster, command, data + +ENDPOINT = 'https://s3.nl-ams.scw.cloud' +BUCKET = 'railiance-platform-pg-backup' +PREFIX = 'platform-pg/application-archives/telemetry/' +LIMIT = 256 * 1024**2 + + +def digest(path): + with path.open('rb') as stream: + return hashlib.file_digest(stream, 'sha256').hexdigest() + + +def run(args, receipt): + source = json.loads(args.source_receipt.read_text()) + if source.get('status') != 'captured' or not source.get('production_resumed'): + raise ValueError('verified_capture_required') + if not 0 < args.source.stat().st_size <= LIMIT or digest(args.source) != source['ciphertext_sha256']: + raise ValueError('source_mismatch') + with args.source.open('rb') as stream: + if stream.read(22) != b'age-encryption.org/v1\n': + raise ValueError('encrypted_artifact_required') + k = ['kubectl', '--kubeconfig', args.kubeconfig] + assert_cluster(k) + values = data(command(k + ['-n', 'databases', 'get', 'secret', 'platform-pg-backup-s3', '-o', 'json']))['data'] + import boto3 + from botocore.config import Config + client = boto3.client('s3', endpoint_url=ENDPOINT, region_name='nl-ams', aws_access_key_id=base64.b64decode(values['ACCESS_KEY_ID']).decode(), aws_secret_access_key=base64.b64decode(values['ACCESS_SECRET_KEY']).decode(), config=Config(signature_version='s3v4', connect_timeout=15, read_timeout=60, retries={'max_attempts': 2}, request_checksum_calculation='when_required', response_checksum_validation='when_required', s3={'addressing_style': 'path'})) + key = PREFIX + datetime.now(timezone.utc).strftime('%Y/%m/%d/%H%M%S-') + secrets.token_hex(16) + '.tar.age' + receipt.update(destination='s3://' + BUCKET + '/' + key, ciphertext_sha256=source['ciphertext_sha256'], ciphertext_bytes=args.source.stat().st_size) + with args.source.open('rb') as stream: + uploaded = client.put_object(Bucket=BUCKET, Key=key, Body=stream, ContentLength=args.source.stat().st_size, ContentType='application/octet-stream', Metadata={'sha256': source['ciphertext_sha256'], 'profile': 'telemetry-essentials-v1'}) + version = uploaded.get('VersionId') + receipt['version_id'] = version + if not version or version == 'null': + raise ValueError('version_pin_missing') + response = client.get_object(Bucket=BUCKET, Key=key, VersionId=version) + if response['ContentLength'] != args.source.stat().st_size: + raise ValueError('download_length_mismatch') + with response['Body'] as body, args.output.open('xb') as output: + args.output.chmod(0o600) + total = 0 + while chunk := body.read(1024 * 1024): + total += len(chunk) + if total > LIMIT: + raise ValueError('download_limit_exceeded') + output.write(chunk) + if digest(args.output) != source['ciphertext_sha256']: + raise ValueError('download_digest_mismatch') + receipt.update(status='fetched_pending_restore', version_pinned=True, download_hash_matches=True) + + +def main(): + p = argparse.ArgumentParser(description=__doc__) + p.add_argument('--kubeconfig', required=True) + for name in ['source', 'source-receipt', 'output', 'receipt']: + p.add_argument('--' + name, type=Path, required=True) + args = p.parse_args() + if args.output.exists(): + p.error('output already exists') + fd = os.open(args.receipt, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600) + receipt = {'schema': 'platform.telemetry-primary-archive.v1', 'status': 'failed', 'started_at': datetime.now(timezone.utc).isoformat()} + try: + run(args, receipt) + except Exception as error: + receipt['error_type'] = type(error).__name__ + finally: + receipt['finished_at'] = datetime.now(timezone.utc).isoformat() + with os.fdopen(fd, 'w') as output: + json.dump(receipt, output, indent=2) + return int(receipt['status'] != 'fetched_pending_restore') + + +if __name__ == '__main__': + raise SystemExit(main()) diff --git a/tests/test_telemetry_primary_archive.py b/tests/test_telemetry_primary_archive.py new file mode 100644 index 0000000..e63d4c2 --- /dev/null +++ b/tests/test_telemetry_primary_archive.py @@ -0,0 +1,40 @@ +import importlib.util +import json +from pathlib import Path +from types import SimpleNamespace +import tempfile +import unittest +from unittest.mock import patch +import sys + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / 'scripts')) +import telemetry_primary_archive as archive + + +class PrimaryArchiveGuards(unittest.TestCase): + def test_unresumed_capture_and_changed_bytes_never_request_credentials(self): + with tempfile.TemporaryDirectory() as directory: + source = Path(directory) / 'source.age' + source.write_bytes(b'age-encryption.org/v1\nfixture') + receipt = Path(directory) / 'receipt.json' + args = SimpleNamespace(source=source, source_receipt=receipt) + for metadata in [ + {'status': 'captured', 'production_resumed': False}, + {'status': 'captured', 'production_resumed': True, 'ciphertext_sha256': 'wrong'}, + ]: + receipt.write_text(json.dumps(metadata)) + with patch.object(archive, 'assert_cluster') as access: + with self.assertRaises(ValueError): + archive.run(args, {}) + access.assert_not_called() + + def test_plaintext_never_requests_credentials(self): + with tempfile.TemporaryDirectory() as directory: + source = Path(directory) / 'source.age' + source.write_bytes(b'plaintext even if the extension claims age') + receipt = Path(directory) / 'receipt.json' + receipt.write_text(json.dumps({'status': 'captured', 'production_resumed': True, 'ciphertext_sha256': archive.digest(source)})) + with patch.object(archive, 'assert_cluster') as access: + with self.assertRaises(ValueError): + archive.run(SimpleNamespace(source=source, source_receipt=receipt), {}) + access.assert_not_called()