import sys sys.dont_write_bytecode = True import argparse import contextlib import hmac import json import os import re from db_backend import canonical_postgres_url, is_postgres_url from docker_depth_experiment import ( DOCKER_DEPTH_QUERY_COUNT, _docker_depth_authority, _experiment_identity_matches, _stored_cohort_plan, apply_docker_depth_cohort_manifest, apply_docker_depth_hold_manifest, apply_docker_depth_reactivation_manifest, apply_docker_depth_resolver_disposition_manifest, apply_docker_depth_resolver_refund_manifest, canonical_docker_depth_plan_hash, generate_docker_depth_cohort_manifest, generate_docker_depth_hold_manifest, generate_docker_depth_reactivation_manifest, generate_docker_depth_resolver_disposition_manifest, generate_docker_depth_resolver_refund_manifest, summarize_docker_depth_fresh_coverage, validate_docker_depth_cohort_manifest, validate_docker_depth_config, validate_docker_depth_hold_manifest, validate_docker_depth_reactivation_manifest, validate_docker_depth_resolver_disposition_manifest, validate_docker_depth_resolver_refund_manifest, validate_dockerhub_discovery_policies, ) from migrate_runtime_safety import ( postgres_migration_guard, require_local_sources_stopped, ) from paths import apply_path_config from postgres_runtime import ( canonical_database_url, load_postgres_environment, verify_cluster_identity, ) from runtime_security import ( ClusterAuthorityLock, MAX_EXTENDED_PRIVATE_JSON_BYTES, read_private_json, reject_reparse_components, require_private_directory, require_private_file, preflight_lifecycle_paths, write_private_json_exclusive, ) from scanner_db import ScannerDB APPLICATION_NAME = 'truf-docker-depth-operator' MANIFEST_MAX_BYTES = MAX_EXTENDED_PRIVATE_JSON_BYTES _SHA256_RE = re.compile(r'^[a-f0-9]{64}$') _DISABLED_ACTIONS = frozenset({ 'generate-cohort', 'apply-cohort', 'generate-hold', 'apply-hold', }) _ENABLED_ACTIONS = frozenset({ 'generate-reactivation', 'apply-reactivation', 'generate-resolver-disposition', 'apply-resolver-disposition', 'generate-resolver-refund', 'apply-resolver-refund', }) def load_config(path): """Load path-expanded YAML only; validation and runtime access are separate.""" try: import yaml except ImportError as exc: raise RuntimeError('PyYAML is required') from exc with open(path, 'r', encoding='utf-8') as handle: return apply_path_config(yaml.safe_load(handle) or {}, path) def provenance_policy_sha256(validated): source = validated.normalized_config['sources']['dockerhub'] policies = validate_dockerhub_discovery_policies( source, validated.experiment.queries, ) hashes = {policy['policy_sha256'] for policy in policies} if len(policies) != DOCKER_DEPTH_QUERY_COUNT or len(hashes) != 1: raise RuntimeError('Docker depth provenance policy authority is ambiguous') return next(iter(hashes)) def _action_name(args): for attribute, name in ( ('status', 'status'), ('generate_cohort_manifest', 'generate-cohort'), ('apply_cohort_manifest', 'apply-cohort'), ('generate_hold_manifest', 'generate-hold'), ('apply_hold_manifest', 'apply-hold'), ('generate_reactivation_manifest', 'generate-reactivation'), ('apply_reactivation_manifest', 'apply-reactivation'), ('generate_resolver_disposition_manifest', 'generate-resolver-disposition'), ('apply_resolver_disposition_manifest', 'apply-resolver-disposition'), ('generate_resolver_refund_manifest', 'generate-resolver-refund'), ('apply_resolver_refund_manifest', 'apply-resolver-refund'), ): if getattr(args, attribute, None): return name raise RuntimeError('Docker depth operator action is unavailable') def _require_action_arguments(parser, args, action): applying = action.startswith('apply-') supplied_apply_option = bool( args.confirm_apply or args.sources_stopped or args.approve_sha256 ) if applying: if not args.confirm_apply or not args.sources_stopped or not args.approve_sha256: parser.error( 'apply actions require --approve-sha256, --apply, and --sources-stopped' ) if not _SHA256_RE.fullmatch(args.approve_sha256): parser.error('--approve-sha256 must be one lowercase SHA-256 value') elif supplied_apply_option: parser.error('approval options are valid only for apply actions') def _require_action_config_state(experiment, action): if action in _DISABLED_ACTIONS and experiment.enabled: raise RuntimeError('Docker depth reviewed preparation requires disabled config') if action in _ENABLED_ACTIONS and not experiment.enabled: raise RuntimeError('Docker depth reviewed release requires enabled config') def _manifest_path(path, *, existing): absolute = reject_reparse_components(os.path.abspath(os.fspath(path))) require_private_directory(os.path.dirname(absolute), create=False) if existing: require_private_file(absolute) elif os.path.lexists(absolute): require_private_file(absolute) return absolute def _publish_manifest(path, manifest): absolute = _manifest_path(path, existing=False) if os.path.lexists(absolute): if read_private_json(absolute, max_bytes=MANIFEST_MAX_BYTES) != manifest: raise RuntimeError('A different reviewed manifest already exists') return absolute, False try: write_private_json_exclusive( absolute, manifest, max_bytes=MANIFEST_MAX_BYTES, ) except FileExistsError: if read_private_json(absolute, max_bytes=MANIFEST_MAX_BYTES) != manifest: raise RuntimeError('Reviewed manifest publication raced a different file') return absolute, False require_private_file(absolute) return absolute, True def _read_approved_manifest(path, validator, experiment, policy_sha256, approved): absolute = _manifest_path(path, existing=True) manifest = read_private_json(absolute, max_bytes=MANIFEST_MAX_BYTES) normalized, manifest_sha256 = validator( manifest, experiment, policy_sha256, ) if not hmac.compare_digest(manifest_sha256, approved): raise ValueError('Reviewed manifest approval hash conflicts') return absolute, normalized, manifest_sha256 def _prepare_action(args, action, experiment, policy_sha256): if action.startswith('generate-'): attribute = action.replace('-', '_') + '_manifest' path = _manifest_path(getattr(args, attribute), existing=False) return {'path': path} if not action.startswith('apply-'): return {} kind = action.removeprefix('apply-') validator = { 'cohort': validate_docker_depth_cohort_manifest, 'hold': validate_docker_depth_hold_manifest, 'reactivation': validate_docker_depth_reactivation_manifest, 'resolver-disposition': validate_docker_depth_resolver_disposition_manifest, 'resolver-refund': validate_docker_depth_resolver_refund_manifest, }[kind] path = getattr(args, f'apply_{kind.replace("-", "_")}_manifest') absolute, manifest, manifest_sha256 = _read_approved_manifest( path, validator, experiment, policy_sha256, args.approve_sha256, ) return { 'path': absolute, 'manifest': manifest, 'manifest_sha256': manifest_sha256, } def _verify_online_cluster_identity(db, dsn, identity): canonical = canonical_postgres_url( dsn, identity['database'], identity['user'], identity['port'], ) if canonical != dsn: raise RuntimeError('Managed PostgreSQL DSN is not canonical') row = db.conn.execute( '''SELECT pg_catalog.current_database() AS database, CURRENT_USER AS user_name, pg_catalog.current_setting('data_directory') AS data_directory, pg_catalog.current_setting('port')::integer AS port, (SELECT system_identifier::text FROM pg_catalog.pg_control_system()) AS system_identifier''' ).fetchone() checks = { 'database': (str(row['database']), str(identity['database'])), 'user': (str(row['user_name']), str(identity['user'])), 'data_directory': ( os.path.normcase(os.path.realpath(os.path.abspath(row['data_directory']))), os.path.normcase(os.path.realpath(os.path.abspath(identity['data_directory']))), ), 'port': (int(row['port']), int(identity['port'])), 'system_identifier': ( str(row['system_identifier']), str(identity['system_identifier']), ), } if any(actual != expected for actual, expected in checks.values()): raise RuntimeError('Online PostgreSQL identity does not match private authority') db.conn.commit() @contextlib.contextmanager def operator_database(config_path, config, *, read_only): preflight_lifecycle_paths(config_path, config) load_postgres_environment(config_path, config) dsn = canonical_database_url() if not dsn or not is_postgres_url(dsn): raise RuntimeError('Canonical managed PostgreSQL DSN is unavailable') with ClusterAuthorityLock(config, endpoint_dsn=dsn): require_local_sources_stopped(config) identity = verify_cluster_identity(config) dsn = canonical_postgres_url( dsn, identity['database'], identity['user'], identity['port'], ) db = ScannerDB(db_url=dsn, initialize=False) try: if not db.enabled: raise RuntimeError('Managed PostgreSQL connection is unavailable') _verify_online_cluster_identity(db, dsn, identity) db.set_application_name(APPLICATION_NAME) with postgres_migration_guard(db): db.require_runtime_safety_schema() db.require_final_cutover() if read_only: db.conn.execute('SET default_transaction_read_only = on') db.conn.commit() yield db finally: db.close() def _status(db, experiment, policy_sha256): authority = _docker_depth_authority(experiment, policy_sha256) coverage = summarize_docker_depth_fresh_coverage( db, experiment, policy_sha256, ) row = db.conn.execute( 'SELECT * FROM docker_depth_experiments WHERE experiment_key = ?', (authority['experiment_key'],), ).fetchone() state = 'absent' plan_sha256 = '' hold_manifest_sha256 = '' fence_active = 0 counts = { 'owned_policy_events': 0, 'planned_queries': 0, 'planned_repositories': 0, 'targets': 0, 'unreleased_holds': 0, } if row: if not _experiment_identity_matches(row, authority): raise RuntimeError('Docker depth persisted authority drifted') state = str(row['state']) plan_sha256 = str(row['plan_sha256'] or '') hold_manifest_sha256 = str(row['hold_manifest_sha256'] or '') fence_active = int(any( row[name] is not None for name in ('fence_owner', 'fence_token', 'fence_expires_at') )) if plan_sha256: stored = _stored_cohort_plan(db.conn, row, authority) if canonical_docker_depth_plan_hash(stored) != plan_sha256: raise RuntimeError('Docker depth persisted cohort hash drifted') count_row = db.conn.execute( '''SELECT (SELECT COUNT(*) FROM docker_depth_experiment_queries WHERE experiment_id = ?) AS planned_queries, (SELECT COUNT(*) FROM docker_depth_experiment_repositories WHERE experiment_id = ?) AS planned_repositories, (SELECT COUNT(*) FROM docker_depth_experiment_targets WHERE experiment_id = ?) AS targets, (SELECT COUNT(*) FROM target_queue_policy_events WHERE experiment_id = ?) AS owned_policy_events, (SELECT COUNT(*) FROM target_queue_policy_events cold_event LEFT JOIN target_queue_policy_events reverse_event ON reverse_event.reverses_event_id = cold_event.id WHERE cold_event.experiment_id = ? AND cold_event.action = 'cold' AND reverse_event.id IS NULL) AS unreleased_holds''', (row['id'], row['id'], row['id'], row['id'], row['id']), ).fetchone() counts = {name: int(count_row[name]) for name in counts} db.conn.commit() return { 'action': 'status', 'config_enabled': bool(experiment.enabled), 'config_sha256': authority['config_sha256'], 'counts': counts, 'experiment_present': bool(row), 'fence_active': fence_active, 'fresh_coverage': coverage, 'hold_manifest_sha256': hold_manifest_sha256, 'ordered_queries_sha256': authority['ordered_queries_sha256'], 'plan_sha256': plan_sha256, 'provenance_policy_sha256': authority['provenance_policy_sha256'], 'selector_sha256': authority['selector_sha256'], 'state': state, } def _execute_action(db, action, prepared, experiment, policy_sha256): if action == 'status': return _status(db, experiment, policy_sha256) if action == 'generate-cohort': manifest, manifest_sha256 = generate_docker_depth_cohort_manifest( db, experiment, policy_sha256, ) path, created = _publish_manifest(prepared['path'], manifest) return { 'action': action, 'files_created': int(created), 'manifest_sha256': manifest_sha256, 'path': os.path.basename(path), 'plan_sha256': manifest['plan_sha256'], 'queries': len(manifest['plan']['queries']), 'repositories': sum( len(item['repositories']) for item in manifest['plan']['queries'] ), } if action == 'apply-cohort': result = apply_docker_depth_cohort_manifest( db, experiment, policy_sha256, prepared['manifest'], prepared['manifest_sha256'], ) return { 'action': action, 'manifest_sha256': prepared['manifest_sha256'], 'path': os.path.basename(prepared['path']), 'plan_sha256': result['plan_sha256'], 'plans_persisted': int(bool(result['planned'])), 'queries': len(result['plan']['queries']), 'repositories': sum( len(item['repositories']) for item in result['plan']['queries'] ), } if action == 'generate-hold': manifest, manifest_sha256 = generate_docker_depth_hold_manifest( db, experiment, policy_sha256, ) path, created = _publish_manifest(prepared['path'], manifest) return { 'action': action, 'conflicts': manifest['conflict_count'], 'entries': manifest['entry_count'], 'files_created': int(created), 'manifest_sha256': manifest_sha256, 'path': os.path.basename(path), 'plan_sha256': manifest['plan_sha256'], } if action == 'apply-hold': result = apply_docker_depth_hold_manifest( db, experiment, policy_sha256, prepared['manifest'], prepared['manifest_sha256'], ) return { 'action': action, 'conflicts': int(result['conflicts']), 'duplicates': int(result['duplicates']), 'manifest_sha256': prepared['manifest_sha256'], 'path': os.path.basename(prepared['path']), 'transitioned': int(result['transitioned']), } if action == 'generate-reactivation': manifest, manifest_sha256 = generate_docker_depth_reactivation_manifest( db, experiment, policy_sha256, ) path, created = _publish_manifest(prepared['path'], manifest) return { 'action': action, 'entries': manifest['entry_count'], 'files_created': int(created), 'hold_manifest_sha256': manifest['hold_manifest_sha256'], 'manifest_sha256': manifest_sha256, 'path': os.path.basename(path), } if action == 'apply-reactivation': result = apply_docker_depth_reactivation_manifest( db, experiment, policy_sha256, prepared['manifest'], prepared['manifest_sha256'], ) return { 'action': action, 'duplicates': int(result['duplicates']), 'manifest_sha256': prepared['manifest_sha256'], 'path': os.path.basename(prepared['path']), 'transitioned': int(result['transitioned']), } if action == 'generate-resolver-disposition': manifest, manifest_sha256 = ( generate_docker_depth_resolver_disposition_manifest( db, experiment, policy_sha256, ) ) path, created = _publish_manifest(prepared['path'], manifest) return { 'action': action, 'files_created': int(created), 'manifest_sha256': manifest_sha256, 'outcome': manifest['entries'][0]['outcome'], 'path': os.path.basename(path), } if action == 'apply-resolver-disposition': result = apply_docker_depth_resolver_disposition_manifest( db, experiment, policy_sha256, prepared['manifest'], prepared['manifest_sha256'], ) return { 'action': action, 'applied': int(result['applied']), 'duplicates': int(result['duplicates']), 'manifest_sha256': prepared['manifest_sha256'], 'outcome': result['outcome'], 'path': os.path.basename(prepared['path']), 'state': result['state'], } if action == 'generate-resolver-refund': manifest, manifest_sha256 = generate_docker_depth_resolver_refund_manifest( db, experiment, policy_sha256, prepared['evidence_log_path'], ) path, created = _publish_manifest(prepared['path'], manifest) return { 'action': action, 'entries': manifest['entry_count'], 'files_created': int(created), 'held_entries': manifest['held_entry_count'], 'manifest_sha256': manifest_sha256, 'path': os.path.basename(path), 'refund_attempts': manifest['refund_attempts'], } if action == 'apply-resolver-refund': result = apply_docker_depth_resolver_refund_manifest( db, experiment, policy_sha256, prepared['manifest'], prepared['manifest_sha256'], prepared['evidence_log_path'], ) return { 'action': action, 'duplicates': int(result['duplicates']), 'manifest_sha256': prepared['manifest_sha256'], 'path': os.path.basename(prepared['path']), 'refunded': int(result['refunded']), 'state': result['state'], } raise RuntimeError('Docker depth operator action is unsupported') def parse_args(argv=None): parser = argparse.ArgumentParser( description='Offline reviewed operator for the bounded Docker depth experiment.', ) parser.add_argument( '--config', default=os.path.join(os.path.dirname(__file__), 'config.yaml'), ) actions = parser.add_mutually_exclusive_group(required=True) actions.add_argument('--status', action='store_true') actions.add_argument('--generate-cohort-manifest') actions.add_argument('--apply-cohort-manifest') actions.add_argument('--generate-hold-manifest') actions.add_argument('--apply-hold-manifest') actions.add_argument('--generate-reactivation-manifest') actions.add_argument('--apply-reactivation-manifest') actions.add_argument('--generate-resolver-disposition-manifest') actions.add_argument('--apply-resolver-disposition-manifest') actions.add_argument('--generate-resolver-refund-manifest') actions.add_argument('--apply-resolver-refund-manifest') parser.add_argument('--approve-sha256') parser.add_argument('--apply', dest='confirm_apply', action='store_true') parser.add_argument('--sources-stopped', action='store_true') args = parser.parse_args(argv) action = _action_name(args) _require_action_arguments(parser, args, action) return args, action def main(argv=None): try: args, action = parse_args(argv) config_path = os.path.abspath(args.config) config = load_config(config_path) validated = validate_docker_depth_config( config, managed_postgres=True, final_cutover=True, ) if validated.experiment is None: raise RuntimeError('Docker depth experiment configuration is unavailable') experiment = validated.experiment _require_action_config_state(experiment, action) policy_sha256 = provenance_policy_sha256(validated) prepared = _prepare_action( args, action, experiment, policy_sha256, ) if action in ('generate-resolver-refund', 'apply-resolver-refund'): log_path = reject_reparse_components(os.path.abspath(os.path.join( validated.normalized_config['global']['log_dir'], 'dockerhub.log', ))) require_private_file(log_path) prepared['evidence_log_path'] = log_path with operator_database( config_path, validated.normalized_config, read_only=not action.startswith('apply-'), ) as db: report = _execute_action( db, action, prepared, experiment, policy_sha256, ) print(json.dumps(report, ensure_ascii=True, sort_keys=True)) return 0 except Exception as exc: raise SystemExit( f'Docker depth operator failed closed: {type(exc).__name__}' ) from None if __name__ == '__main__': main()