Files
truf-server/app/docker_depth_operator.py
2026-09-30 20:30:56 +03:00

547 lines
22 KiB
Python

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()