from dataclasses import dataclass, field import hashlib import hmac import json import os from pathlib import Path import re import stat from runtime_document import ( MAX_CONFIG_DOCUMENT_BYTES, MAX_SECRETS_DOCUMENT_BYTES, RuntimeDocumentError, _resolve_package_manifest_path, load_yaml_document, preview_runtime_documents, ) from runtime_security import ( PrivateFileLock, durable_replace, fsync_directory, harden_private_file, read_stable_root_file, require_private_directory, require_private_file, ) from worker_package import ( MAX_WORKER_PACKAGE_MANIFEST_BYTES, load_worker_package_manifest_bytes, ) MANAGED_TEMPLATE_PATH = Path(__file__).with_name('config.linux.yaml') MANAGED_WORKER_PACKAGE_DIRECTORY = Path('/data/worker-packages') MANAGED_SECRETS_PATH = Path('/data/config/secrets.yaml') MANAGED_CANDIDATE_DIRECTORY = Path('/data/runtime-document-candidates') MANAGED_CONFIG_CANDIDATE_PATH = MANAGED_CANDIDATE_DIRECTORY / 'config.yaml' MANAGED_SECRETS_CANDIDATE_PATH = MANAGED_CANDIDATE_DIRECTORY / 'secrets.yaml' MANAGED_CANDIDATE_LOCK_PATH = MANAGED_CANDIDATE_DIRECTORY / 'candidate.lock' MAX_CONFIG_DIFF_ENTRIES = 200 MAX_CONFIG_DIFF_OUTPUT_BYTES = 65536 MAX_CONFIG_DIFF_PATH_BYTES = 512 MAX_CONFIG_DIFF_VALUE_BYTES = 512 MAX_CONFIG_DIFF_ENTRY_BYTES = 1024 _HASH = re.compile(r'^[0-9a-f]{64}$') _CANDIDATE_TEMPORARY = re.compile( r'^\.(?:config|secrets)\.yaml\.[0-9a-f]{24}\.tmp$' ) _SENSITIVE_PATH = re.compile( r'(?:pass(?:word)?|token|secret|credential|auth|cookie)', re.IGNORECASE, ) _SENSITIVE_CONFIG_FIELDS = frozenset({ 'dashboard_db_url', 'database_url', 'edge_marker', 'postgres_dsn', }) _MISSING = object() @dataclass(frozen=True) class ManagedRuntimeDocuments: config: dict secrets: dict | None config_sha256: str secrets_sha256: str | None @dataclass(frozen=True) class RuntimeDocumentIdentity: sha256: str byte_count: int present: bool @dataclass(frozen=True) class ManagedRuntimeDocumentState: active_config: RuntimeDocumentIdentity active_secrets: RuntimeDocumentIdentity candidate_config: RuntimeDocumentIdentity candidate_secrets: RuntimeDocumentIdentity @dataclass(frozen=True) class ConfigDiffEntry: path: str change: str before: str | None after: str | None value_redacted: bool value_truncated: bool @dataclass(frozen=True) class ConfigDiff: entries: tuple[ConfigDiffEntry, ...] truncated: bool format_only_changed: bool output_bytes: int @dataclass(frozen=True) class SecretsRedactedDiff: document_changed: bool semantic_changed: bool pools_before: int pools_after: int pools_added: int pools_removed: int pools_changed: int entries_before: int entries_after: int entries_added: int entries_removed: int pools_reordered: int usernames_added: int usernames_removed: int usernames_changed: int tokens_changed: int @dataclass(frozen=True) class ManagedRuntimeCandidatePreview: document: str state: ManagedRuntimeDocumentState proposed: RuntimeDocumentIdentity diff: ConfigDiff | SecretsRedactedDiff @dataclass(frozen=True) class ManagedRuntimeCandidateRevision: document: str before: ManagedRuntimeDocumentState after: ManagedRuntimeDocumentState proposed: RuntimeDocumentIdentity created: bool content_changed: bool written: bool diff: ConfigDiff | SecretsRedactedDiff @dataclass(frozen=True) class ManagedRuntimeApplyVerification: action: str state: ManagedRuntimeDocumentState config_source: str secrets_source: str effective_config_sha256: str effective_secrets_sha256: str @dataclass(frozen=True) class ManagedRuntimeEditorDocument: document: str source: str text: str = field(repr=False) state: ManagedRuntimeDocumentState selected: RuntimeDocumentIdentity @dataclass(frozen=True) class _CandidateSnapshot: active_config: bytes active_secrets: bytes candidate_config: bytes | None candidate_secrets: bytes | None state: ManagedRuntimeDocumentState def _error_details(exc, default_document=None): return ( exc.category, exc.line, exc.column, exc.document or default_document, exc.path, ) def _read_private_document(path, max_bytes, document): failed = False payload = None descriptor = None try: checked = require_private_file(os.fspath(path)) flags = os.O_RDONLY if hasattr(os, 'O_BINARY'): flags |= os.O_BINARY if hasattr(os, 'O_NOFOLLOW'): flags |= os.O_NOFOLLOW descriptor = os.open(checked, flags) before = os.fstat(descriptor) if not stat.S_ISREG(before.st_mode) or before.st_nlink != 1: raise OSError('private document identity is invalid') with os.fdopen(descriptor, 'rb') as handle: descriptor = None payload = handle.read(max_bytes + 1) after = os.fstat(handle.fileno()) require_private_file(checked) current = os.stat(checked, follow_symlinks=False) identity = lambda value: (value.st_dev, value.st_ino) if ( identity(before) != identity(after) or identity(after) != identity(current) or after.st_nlink != 1 or current.st_nlink != 1 or before.st_size != after.st_size or after.st_size != current.st_size or getattr(before, 'st_mtime_ns', None) != getattr(after, 'st_mtime_ns', None) or getattr(after, 'st_mtime_ns', None) != getattr(current, 'st_mtime_ns', None) or getattr(before, 'st_ctime_ns', None) != getattr(after, 'st_ctime_ns', None) or ( os.name != 'nt' and getattr(after, 'st_ctime_ns', None) != getattr(current, 'st_ctime_ns', None) ) ): raise OSError('private document changed while it was being read') except Exception: failed = True except BaseException: payload = None raise finally: if descriptor is not None: os.close(descriptor) if failed: return None, ('schema', None, None, document, 'root') if len(payload) > max_bytes: return None, ('size', None, None, document, None) return payload, None def _read_package_manifest(path): if os.name == 'nt': return _read_private_document( path, MAX_WORKER_PACKAGE_MANIFEST_BYTES, 'package_capabilities', ) try: return read_stable_root_file( path, MAX_WORKER_PACKAGE_MANIFEST_BYTES, MANAGED_WORKER_PACKAGE_DIRECTORY, ), None except Exception: return None, ('schema', None, None, 'package_capabilities', 'root') def _package_capability_evidence(config): profiles = evidence = profile_name = profile = reference = None resolved = payload = failure = manifest = None try: try: profiles = config['supervisor']['worker_api']['compatibility_profiles'] except Exception: return None, None if type(profiles) is not dict: return None, None evidence = {} for profile_name, profile in profiles.items(): if type(profile) is not dict: return None, None reference = profile.get('package_manifest') resolved = _resolve_package_manifest_path(config, reference) if resolved is None: return None, None try: payload, failure = _read_package_manifest(resolved) if failure is not None: return None, None manifest = load_worker_package_manifest_bytes(payload) except Exception: return None, None evidence[profile_name] = { 'package_manifest': reference, 'capabilities': manifest['capabilities'], } return evidence, True finally: config = profiles = evidence = profile_name = profile = reference = None resolved = payload = failure = manifest = None def _validate_managed_runtime_files_inner(config_path, secrets_bytes): config_payload = None template_payload = None active_secrets_payload = None config = None package_capabilities = None validated = None try: config_payload, failure = _read_private_document( config_path, MAX_CONFIG_DOCUMENT_BYTES, 'config', ) if failure is not None: return None, failure template_payload, failure = _read_private_document( MANAGED_TEMPLATE_PATH, MAX_CONFIG_DOCUMENT_BYTES, 'schema', ) if failure is not None: return None, failure if secrets_bytes is None: active_secrets_payload, failure = _read_private_document( MANAGED_SECRETS_PATH, MAX_SECRETS_DOCUMENT_BYTES, 'secrets', ) if failure is not None: return None, failure else: active_secrets_payload = secrets_bytes config = load_yaml_document( config_payload, max_bytes=MAX_CONFIG_DOCUMENT_BYTES, ) package_capabilities, ready = _package_capability_evidence(config) if ready is None: return None, ('capability', None, None, 'package_capabilities', 'root') validated = preview_runtime_documents( config_payload, active_secrets_payload, config_template_payload=template_payload, package_capabilities=package_capabilities, ) return ManagedRuntimeDocuments( config=validated.config, secrets=validated.secrets, config_sha256=hashlib.sha256(config_payload).hexdigest(), secrets_sha256=hashlib.sha256(active_secrets_payload).hexdigest(), ), None except RuntimeDocumentError as exc: return None, _error_details(exc, 'config') except Exception: return None, ('schema', None, None, 'schema', 'root') finally: secrets_bytes = None config_payload = None template_payload = None active_secrets_payload = None config = None package_capabilities = None validated = None def validate_managed_runtime_files(config_path, *, secrets_bytes=None): validated = failure = None try: validated, failure = _validate_managed_runtime_files_inner( config_path, secrets_bytes, ) finally: config_path = None secrets_bytes = None if failure is not None: raise RuntimeDocumentError( failure[0], failure[1], failure[2], document=failure[3], path=failure[4], ) return validated def _load_managed_runtime_config_inner(config_path): payload = config = failure = None try: payload, failure = _read_private_document( config_path, MAX_CONFIG_DOCUMENT_BYTES, 'config', ) if failure is not None: return None, failure config = load_yaml_document(payload, max_bytes=MAX_CONFIG_DOCUMENT_BYTES) if type(config) is not dict: return None, ('mapping_root', None, None, 'config', 'root') return ManagedRuntimeDocuments( config=config, secrets=None, config_sha256=hashlib.sha256(payload).hexdigest(), secrets_sha256=None, ), None except RuntimeDocumentError as exc: return None, _error_details(exc, 'config') except Exception: return None, ('schema', None, None, 'config', 'root') finally: payload = config = failure = None def load_managed_runtime_config(config_path): loaded, failure = _load_managed_runtime_config_inner(config_path) config_path = None if failure is not None: raise RuntimeDocumentError( failure[0], failure[1], failure[2], document=failure[3], path=failure[4], ) return loaded def _candidate_failure(category='schema', document='schema', path='root'): return category, None, None, document, path def _document_identity(payload, present=True): try: return RuntimeDocumentIdentity( sha256=hashlib.sha256(payload).hexdigest(), byte_count=len(payload), present=present, ) finally: payload = None def _read_optional_candidate(path, max_bytes, document): if not os.path.lexists(path): return None, None return _read_private_document(path, max_bytes, document) def _read_candidate_snapshot(config_path): active_config = active_secrets = None candidate_config = candidate_secrets = None try: active_config, failure = _read_private_document( config_path, MAX_CONFIG_DOCUMENT_BYTES, 'config', ) if failure is not None: return None, failure active_secrets, failure = _read_private_document( MANAGED_SECRETS_PATH, MAX_SECRETS_DOCUMENT_BYTES, 'secrets', ) if failure is not None: return None, failure candidate_config, failure = _read_optional_candidate( MANAGED_CONFIG_CANDIDATE_PATH, MAX_CONFIG_DOCUMENT_BYTES, 'config', ) if failure is not None: return None, failure candidate_secrets, failure = _read_optional_candidate( MANAGED_SECRETS_CANDIDATE_PATH, MAX_SECRETS_DOCUMENT_BYTES, 'secrets', ) if failure is not None: return None, failure state = ManagedRuntimeDocumentState( active_config=_document_identity(active_config), active_secrets=_document_identity(active_secrets), candidate_config=_document_identity( candidate_config if candidate_config is not None else active_config, candidate_config is not None, ), candidate_secrets=_document_identity( candidate_secrets if candidate_secrets is not None else active_secrets, candidate_secrets is not None, ), ) return _CandidateSnapshot( active_config=active_config, active_secrets=active_secrets, candidate_config=candidate_config, candidate_secrets=candidate_secrets, state=state, ), None except Exception: return None, _candidate_failure() except BaseException: active_config = active_secrets = None candidate_config = candidate_secrets = None raise def _validate_payload_pair(config_payload, secrets_payload): template_payload = config = package_capabilities = validated = None try: template_payload, failure = _read_private_document( MANAGED_TEMPLATE_PATH, MAX_CONFIG_DOCUMENT_BYTES, 'schema', ) if failure is not None: return None, failure config = load_yaml_document( config_payload, max_bytes=MAX_CONFIG_DOCUMENT_BYTES, ) package_capabilities, ready = _package_capability_evidence(config) if ready is None: return None, _candidate_failure( 'capability', 'package_capabilities', 'root', ) validated = preview_runtime_documents( config_payload, secrets_payload, config_template_payload=template_payload, package_capabilities=package_capabilities, ) return validated, None except RuntimeDocumentError as exc: return None, _error_details(exc, 'config') except Exception: return None, _candidate_failure() finally: config_payload = None secrets_payload = None template_payload = None config = None package_capabilities = None validated = None def _prepare_candidate_store(): existed = os.path.lexists(MANAGED_CANDIDATE_DIRECTORY) require_private_directory(MANAGED_CANDIDATE_DIRECTORY, create=True) if not existed: fsync_directory(Path(MANAGED_CANDIDATE_DIRECTORY).parent) def _verify_candidate_lock(): checked = require_private_file(MANAGED_CANDIDATE_LOCK_PATH) details = os.stat(checked, follow_symlinks=False) if not stat.S_ISREG(details.st_mode) or details.st_nlink != 1: raise OSError('candidate lock identity is invalid') if os.name != 'nt' and stat.S_IMODE(details.st_mode) != 0o600: raise OSError('candidate lock mode is invalid') def _cleanup_candidate_temporaries(): removed = False with os.scandir(MANAGED_CANDIDATE_DIRECTORY) as entries: for entry in entries: if _CANDIDATE_TEMPORARY.fullmatch(entry.name) is None: continue details = os.stat(entry.path, follow_symlinks=False) if ( not stat.S_ISREG(details.st_mode) or details.st_nlink != 1 or (os.name != 'nt' and stat.S_IMODE(details.st_mode) != 0o600) ): raise OSError('candidate temporary file identity is invalid') require_private_file(entry.path) os.unlink(entry.path) removed = True if removed: fsync_directory(MANAGED_CANDIDATE_DIRECTORY) def _valid_expected_hash(value): return type(value) is str and _HASH.fullmatch(value) is not None def _same_hash(actual, expected): return hmac.compare_digest(actual, expected) def _truncate_text(value, max_bytes): encoded = None shortened = None try: encoded = value.encode('utf-8') if len(encoded) <= max_bytes: return value, False shortened = encoded[:max(0, max_bytes - 3)].decode('utf-8', 'ignore') + '...' return shortened, True finally: value = None encoded = None shortened = None def _render_diff_value(value): try: return json.dumps( value, ensure_ascii=True, sort_keys=True, separators=(',', ':'), ) finally: value = None def _config_diff(before, after, before_hash, after_hash): entries = [] output_bytes = 2 # JSON list delimiters. truncated = False stack = [('', before, after)] old = new = old_text = new_text = entry = None try: while stack and not truncated: path, old, new = stack.pop() if old == new: continue old_empty_container = type(old) in (dict, list) and not old new_empty_container = type(new) in (dict, list) and not new leaf_container_change = ( old_empty_container and (new is _MISSING or type(new) is not type(old)) ) or ( new_empty_container and (old is _MISSING or type(old) is not type(new)) ) if not leaf_container_change and type(old) is dict and type(new) is dict: keys = sorted(set(old) | set(new), reverse=True) for key in keys: child = f'{path}.{key}' if path else str(key) stack.append(( child, old.get(key, _MISSING), new.get(key, _MISSING), )) continue if not leaf_container_change and type(old) is list and type(new) is list: for index in range(max(len(old), len(new)) - 1, -1, -1): stack.append(( f'{path}[{index}]', old[index] if index < len(old) else _MISSING, new[index] if index < len(new) else _MISSING, )) continue if ( not leaf_container_change and old is _MISSING and type(new) in (dict, list) and new ): stack.append((path, {} if type(new) is dict else [], new)) continue if ( not leaf_container_change and new is _MISSING and type(old) in (dict, list) and old ): stack.append((path, old, {} if type(old) is dict else [])) continue if not leaf_container_change and ( type(old) in (dict, list) or type(new) in (dict, list) ): if new is not _MISSING: stack.append(( path, {} if type(new) is dict else [] if type(new) is list else _MISSING, new, )) if old is not _MISSING: stack.append(( path, old, {} if type(old) is dict else [] if type(old) is list else _MISSING, )) continue path = path or 'root' if len(entries) >= MAX_CONFIG_DIFF_ENTRIES: truncated = True continue if len(path.encode('utf-8')) > MAX_CONFIG_DIFF_PATH_BYTES: truncated = True continue lowered = '.' + path.lower() + '.' field_name = re.split(r'[.\[]', path)[-1].rstrip(']').lower() redacted = ( '.env.' in lowered or _SENSITIVE_PATH.search(path) is not None or field_name in _SENSITIVE_CONFIG_FIELDS or type(old) is str or type(new) is str ) value_truncated = False if redacted: old_text = None if old is _MISSING else '[redacted]' new_text = None if new is _MISSING else '[redacted]' else: old_text = None if old is _MISSING else _render_diff_value(old) new_text = None if new is _MISSING else _render_diff_value(new) if old_text is not None: old_text, cut = _truncate_text( old_text, MAX_CONFIG_DIFF_VALUE_BYTES, ) value_truncated = value_truncated or cut if new_text is not None: new_text, cut = _truncate_text( new_text, MAX_CONFIG_DIFF_VALUE_BYTES, ) value_truncated = value_truncated or cut change = ( 'added' if old is _MISSING else 'removed' if new is _MISSING else 'changed' ) entry = ConfigDiffEntry( path=path, change=change, before=old_text, after=new_text, value_redacted=redacted, value_truncated=value_truncated, ) entry_bytes = len(json.dumps( entry.__dict__, ensure_ascii=True, sort_keys=True, separators=(',', ':'), ).encode('utf-8')) entry_output_bytes = entry_bytes + (1 if entries else 0) if ( entry_bytes > MAX_CONFIG_DIFF_ENTRY_BYTES or output_bytes + entry_output_bytes > MAX_CONFIG_DIFF_OUTPUT_BYTES ): truncated = True continue entries.append(entry) output_bytes += entry_output_bytes return ConfigDiff( entries=tuple(entries), truncated=truncated, format_only_changed=( before_hash != after_hash and before == after and not entries ), output_bytes=output_bytes, ) finally: before = after = old = new = None old_text = new_text = entry = None stack = None entries = None def _secrets_diff(before, after, before_hash, after_hash): before_pools = after_pools = None old_entries = new_entries = None old_username = new_username = None entry = entries = None before_names = after_names = common_pools = None old_names = new_names = None pool_name = name = None try: before_pools = before.get('auth_pools', {}) after_pools = after.get('auth_pools', {}) before_names = set(before_pools) after_names = set(after_pools) common_pools = before_names & after_names entries_before = entries_after = 0 for entries in before_pools.values(): entries_before += len(entries) for entries in after_pools.values(): entries_after += len(entries) entries_added = entries_removed = 0 usernames_added = usernames_removed = usernames_changed = 0 tokens_changed = 0 for pool_name in before_names | after_names: old_entries = {} for entry in before_pools.get(pool_name, []): old_entries[entry['name']] = entry new_entries = {} for entry in after_pools.get(pool_name, []): new_entries[entry['name']] = entry old_names = set(old_entries) new_names = set(new_entries) entries_added += len(new_names - old_names) entries_removed += len(old_names - new_names) for name in new_names - old_names: usernames_added += int('username' in new_entries[name]) for name in old_names - new_names: usernames_removed += int('username' in old_entries[name]) for name in old_names & new_names: old_username = old_entries[name].get('username', _MISSING) new_username = new_entries[name].get('username', _MISSING) if old_username is _MISSING and new_username is not _MISSING: usernames_added += 1 elif old_username is not _MISSING and new_username is _MISSING: usernames_removed += 1 elif old_username != new_username: usernames_changed += 1 if old_entries[name].get('token') != new_entries[name].get('token'): tokens_changed += 1 pools_changed = 0 for name in common_pools: pools_changed += int(before_pools[name] != after_pools[name]) return SecretsRedactedDiff( document_changed=before_hash != after_hash, semantic_changed=before != after, pools_before=len(before_pools), pools_after=len(after_pools), pools_added=len(after_names - before_names), pools_removed=len(before_names - after_names), pools_changed=pools_changed, entries_before=entries_before, entries_after=entries_after, entries_added=entries_added, entries_removed=entries_removed, pools_reordered=int( before_names == after_names and list(before_pools) != list(after_pools) ), usernames_added=usernames_added, usernames_removed=usernames_removed, usernames_changed=usernames_changed, tokens_changed=tokens_changed, ) finally: before = after = None before_pools = after_pools = None old_entries = new_entries = None old_username = new_username = entry = entries = None before_names = after_names = common_pools = None old_names = new_names = None pool_name = name = None def _candidate_diff(document, active, proposed, active_payload, proposed_payload): try: active_hash = hashlib.sha256(active_payload).hexdigest() proposed_hash = hashlib.sha256(proposed_payload).hexdigest() if document == 'config': return _config_diff( active.config, proposed.config, active_hash, proposed_hash, ) return _secrets_diff( active.secrets, proposed.secrets, active_hash, proposed_hash, ) finally: active = proposed = None active_payload = proposed_payload = None def _proposed_pair(snapshot, document, candidate_bytes): try: if document == 'config': return ( candidate_bytes, snapshot.candidate_secrets if snapshot.candidate_secrets is not None else snapshot.active_secrets, ) return ( snapshot.candidate_config if snapshot.candidate_config is not None else snapshot.active_config, candidate_bytes, ) finally: snapshot = None candidate_bytes = None def _stage_candidate(path, payload, max_bytes): temporary = None descriptor = None stored = None try: for _attempt in range(16): temporary = Path(path).with_name( f'.{Path(path).name}.{os.urandom(12).hex()}.tmp' ) flags = os.O_WRONLY | os.O_CREAT | os.O_EXCL if hasattr(os, 'O_BINARY'): flags |= os.O_BINARY if hasattr(os, 'O_NOFOLLOW'): flags |= os.O_NOFOLLOW try: descriptor = os.open(temporary, flags, 0o600) break except FileExistsError: temporary = None if descriptor is None or temporary is None: raise OSError('candidate temporary file could not be created') with os.fdopen(descriptor, 'wb') as handle: descriptor = None before = os.fstat(handle.fileno()) if not stat.S_ISREG(before.st_mode) or before.st_nlink != 1: raise OSError('candidate temporary file identity is invalid') handle.write(payload) handle.flush() os.fsync(handle.fileno()) after = os.fstat(handle.fileno()) if ( not stat.S_ISREG(after.st_mode) or after.st_nlink != 1 or after.st_size != len(payload) or (os.name != 'nt' and stat.S_IMODE(after.st_mode) != 0o600) ): raise OSError('candidate temporary file metadata is invalid') harden_private_file(temporary) stored, failure = _read_private_document( path=temporary, max_bytes=max_bytes, document='schema', ) if failure is not None or not hmac.compare_digest( hashlib.sha256(stored).digest(), hashlib.sha256(payload).digest(), ): raise OSError('candidate temporary file verification failed') return temporary except BaseException: if temporary is not None: try: os.unlink(temporary) except OSError: pass raise finally: if descriptor is not None: os.close(descriptor) payload = None stored = None def _rollback_candidate(path, previous_payload, expected_payload, max_bytes, document): current = temporary = None try: current, failure = _read_private_document(path, max_bytes, document) if failure is not None or not hmac.compare_digest( hashlib.sha256(current).digest(), hashlib.sha256(expected_payload).digest(), ): return False if previous_payload is None: os.unlink(path) fsync_directory(MANAGED_CANDIDATE_DIRECTORY) return not os.path.lexists(path) temporary = _stage_candidate(path, previous_payload, max_bytes) durable_replace(temporary, path) temporary = None failure = _verify_final_candidate( path, previous_payload, max_bytes, document, ) if failure is not None: return False fsync_directory(MANAGED_CANDIDATE_DIRECTORY) return True except Exception: return False finally: try: if temporary is not None: try: os.unlink(temporary) except OSError: pass finally: current = temporary = failure = None previous_payload = expected_payload = None def _rollback_committed_candidate( path, previous_payload, expected_payload, max_bytes, document, ): try: with PrivateFileLock(MANAGED_CANDIDATE_LOCK_PATH): _verify_candidate_lock() return _rollback_candidate( path, previous_payload, expected_payload, max_bytes, document, ) except BaseException: return False finally: path = previous_payload = expected_payload = None max_bytes = document = None def _verify_final_candidate(path, expected_payload, max_bytes, document): stored = None try: stored, failure = _read_private_document(path, max_bytes, document) if failure is not None: return failure details = os.stat(path, follow_symlinks=False) if ( not stat.S_ISREG(details.st_mode) or details.st_nlink != 1 or (os.name != 'nt' and stat.S_IMODE(details.st_mode) != 0o600) or len(stored) != len(expected_payload) or not hmac.compare_digest( hashlib.sha256(stored).digest(), hashlib.sha256(expected_payload).digest(), ) ): return _candidate_failure('schema', document, 'root') return None finally: stored = None expected_payload = None def _load_editor_document_inner(config_path, document): snapshot = payload = text = None try: if document not in ('config', 'secrets'): return None, _candidate_failure('invalid_input') _prepare_candidate_store() with PrivateFileLock(MANAGED_CANDIDATE_LOCK_PATH): _verify_candidate_lock() _cleanup_candidate_temporaries() snapshot, failure = _read_candidate_snapshot(config_path) if failure is not None: return None, failure if document == 'config': present = snapshot.state.candidate_config.present payload = snapshot.candidate_config if present else snapshot.active_config selected = snapshot.state.candidate_config else: present = snapshot.state.candidate_secrets.present payload = snapshot.candidate_secrets if present else snapshot.active_secrets selected = snapshot.state.candidate_secrets try: text = payload.decode('utf-8', errors='strict') except UnicodeDecodeError: return None, _candidate_failure('encoding', document, 'root') return ManagedRuntimeEditorDocument( document=document, source='candidate' if present else 'active', text=text, state=snapshot.state, selected=selected, ), None except Exception: return None, _candidate_failure() finally: snapshot = payload = text = None def load_managed_runtime_editor_document(config_path, document): editor = failure = None try: editor, failure = _load_editor_document_inner(config_path, document) finally: config_path = document = None if failure is not None: raise RuntimeDocumentError( failure[0], failure[1], failure[2], document=failure[3], path=failure[4], ) return editor def _preview_candidate_inner(config_path, document, candidate_bytes): snapshot = active_validated = proposed_validated = None proposed_config = proposed_secrets = None active_payload = None try: if document not in ('config', 'secrets'): return None, _candidate_failure('invalid_input') if type(candidate_bytes) is not bytes: return None, _candidate_failure('invalid_input', document, 'root') _prepare_candidate_store() with PrivateFileLock(MANAGED_CANDIDATE_LOCK_PATH): _verify_candidate_lock() _cleanup_candidate_temporaries() snapshot, failure = _read_candidate_snapshot(config_path) if failure is not None: return None, failure active_validated, failure = _validate_payload_pair( snapshot.active_config, snapshot.active_secrets, ) if failure is not None: return None, failure proposed_config, proposed_secrets = _proposed_pair( snapshot, document, candidate_bytes, ) proposed_validated, failure = _validate_payload_pair( proposed_config, proposed_secrets, ) if failure is not None: return None, failure active_payload = ( snapshot.active_config if document == 'config' else snapshot.active_secrets ) return ManagedRuntimeCandidatePreview( document=document, state=snapshot.state, proposed=_document_identity(candidate_bytes), diff=_candidate_diff( document, active_validated, proposed_validated, active_payload, candidate_bytes, ), ), None except Exception: return None, _candidate_failure() finally: candidate_bytes = None snapshot = None active_validated = None proposed_validated = None proposed_config = None proposed_secrets = None active_payload = None def preview_managed_runtime_candidate(config_path, document, candidate_bytes): preview = failure = None try: preview, failure = _preview_candidate_inner( config_path, document, candidate_bytes, ) finally: config_path = None document = None candidate_bytes = None if failure is not None: raise RuntimeDocumentError( failure[0], failure[1], failure[2], document=failure[3], path=failure[4], ) return preview def _save_candidate_inner( config_path, document, candidate_bytes, expected_hashes, ): snapshot = current = after = None active_validated = proposed_validated = None proposed_config = proposed_secrets = None previous_payload = expected_after = None temporary = None committed = False try: if document not in ('config', 'secrets'): return None, _candidate_failure('invalid_input') if type(candidate_bytes) is not bytes: return None, _candidate_failure('invalid_input', document, 'root') if not all(_valid_expected_hash(value) for value in expected_hashes): return None, _candidate_failure('reference', document, 'revision') _prepare_candidate_store() with PrivateFileLock(MANAGED_CANDIDATE_LOCK_PATH): _verify_candidate_lock() _cleanup_candidate_temporaries() snapshot, failure = _read_candidate_snapshot(config_path) if failure is not None: return None, failure actual_hashes = ( snapshot.state.active_config.sha256, snapshot.state.active_secrets.sha256, snapshot.state.candidate_config.sha256, snapshot.state.candidate_secrets.sha256, ) if not all( _same_hash(actual, expected) for actual, expected in zip(actual_hashes, expected_hashes) ): return None, _candidate_failure('reference', document, 'revision') active_validated, failure = _validate_payload_pair( snapshot.active_config, snapshot.active_secrets, ) if failure is not None: return None, failure proposed_config, proposed_secrets = _proposed_pair( snapshot, document, candidate_bytes, ) proposed_validated, failure = _validate_payload_pair( proposed_config, proposed_secrets, ) if failure is not None: return None, failure selected_identity = ( snapshot.state.candidate_config if document == 'config' else snapshot.state.candidate_secrets ) proposed_identity = _document_identity(candidate_bytes) selected_present = selected_identity.present content_changed = not _same_hash( selected_identity.sha256, proposed_identity.sha256, ) or selected_identity.byte_count != proposed_identity.byte_count diff = _candidate_diff( document, active_validated, proposed_validated, snapshot.active_config if document == 'config' else snapshot.active_secrets, candidate_bytes, ) if selected_present and not content_changed: return ManagedRuntimeCandidateRevision( document=document, before=snapshot.state, after=snapshot.state, proposed=proposed_identity, created=False, content_changed=False, written=False, diff=diff, ), None path = ( MANAGED_CONFIG_CANDIDATE_PATH if document == 'config' else MANAGED_SECRETS_CANDIDATE_PATH ) max_bytes = ( MAX_CONFIG_DOCUMENT_BYTES if document == 'config' else MAX_SECRETS_DOCUMENT_BYTES ) previous_payload = ( snapshot.candidate_config if document == 'config' else snapshot.candidate_secrets ) temporary = _stage_candidate(path, candidate_bytes, max_bytes) current, failure = _read_candidate_snapshot(config_path) if failure is not None: return None, failure if current.state != snapshot.state: return None, _candidate_failure('reference', document, 'revision') committed = True durable_replace(temporary, path) temporary = None failure = _verify_final_candidate( path, candidate_bytes, max_bytes, document, ) if failure is not None: rolled_back = _rollback_candidate( path, previous_payload, candidate_bytes, max_bytes, document, ) committed = not rolled_back return None, failure if rolled_back else _candidate_failure() try: fsync_directory(MANAGED_CANDIDATE_DIRECTORY) except Exception: rolled_back = _rollback_candidate( path, previous_payload, candidate_bytes, max_bytes, document, ) committed = not rolled_back return None, _candidate_failure() after, failure = _read_candidate_snapshot(config_path) if failure is not None: rolled_back = _rollback_candidate( path, previous_payload, candidate_bytes, max_bytes, document, ) committed = not rolled_back return None, failure if rolled_back else _candidate_failure() expected_after = ManagedRuntimeDocumentState( active_config=snapshot.state.active_config, active_secrets=snapshot.state.active_secrets, candidate_config=( proposed_identity if document == 'config' else snapshot.state.candidate_config ), candidate_secrets=( proposed_identity if document == 'secrets' else snapshot.state.candidate_secrets ), ) if after.state != expected_after: rolled_back = _rollback_candidate( path, previous_payload, candidate_bytes, max_bytes, document, ) committed = not rolled_back return None, _candidate_failure( 'reference' if rolled_back else 'schema', document, 'revision' if rolled_back else 'root', ) committed = False return ManagedRuntimeCandidateRevision( document=document, before=snapshot.state, after=after.state, proposed=proposed_identity, created=not selected_present, content_changed=content_changed, written=True, diff=diff, ), None except Exception: if committed: _rollback_committed_candidate( path, previous_payload, candidate_bytes, max_bytes, document, ) return None, _candidate_failure() except BaseException: if committed: _rollback_committed_candidate( path, previous_payload, candidate_bytes, max_bytes, document, ) raise finally: if temporary is not None: try: os.unlink(temporary) except OSError: pass candidate_bytes = None snapshot = None current = None after = None active_validated = None proposed_validated = None proposed_config = None proposed_secrets = None previous_payload = None expected_after = None def save_managed_runtime_candidate( config_path, document, candidate_bytes, *, expected_active_config_sha256, expected_active_secrets_sha256, expected_candidate_config_sha256, expected_candidate_secrets_sha256, ): revision = failure = None try: revision, failure = _save_candidate_inner( config_path, document, candidate_bytes, ( expected_active_config_sha256, expected_active_secrets_sha256, expected_candidate_config_sha256, expected_candidate_secrets_sha256, ), ) finally: config_path = None document = None candidate_bytes = None expected_active_config_sha256 = None expected_active_secrets_sha256 = None expected_candidate_config_sha256 = None expected_candidate_secrets_sha256 = None if failure is not None: raise RuntimeDocumentError( failure[0], failure[1], failure[2], document=failure[3], path=failure[4], ) return revision def _verify_candidates_inner( config_path, action, expected_active_config_sha256, expected_active_secrets_sha256, expected_candidate_config_sha256, expected_candidate_secrets_sha256, ): snapshot = validated = None config_payload = secrets_payload = None try: if action not in ('apply-config', 'apply-secrets', 'apply-both'): return None, _candidate_failure('invalid_input') if not _valid_expected_hash(expected_active_config_sha256) or not _valid_expected_hash( expected_active_secrets_sha256 ): return None, _candidate_failure('reference', 'schema', 'revision') required_config = action in ('apply-config', 'apply-both') required_secrets = action in ('apply-secrets', 'apply-both') if ( required_config != (expected_candidate_config_sha256 is not None) or required_secrets != (expected_candidate_secrets_sha256 is not None) or ( expected_candidate_config_sha256 is not None and not _valid_expected_hash(expected_candidate_config_sha256) ) or ( expected_candidate_secrets_sha256 is not None and not _valid_expected_hash(expected_candidate_secrets_sha256) ) ): return None, _candidate_failure('reference', 'schema', 'revision') _prepare_candidate_store() with PrivateFileLock(MANAGED_CANDIDATE_LOCK_PATH): _verify_candidate_lock() _cleanup_candidate_temporaries() snapshot, failure = _read_candidate_snapshot(config_path) if failure is not None: return None, failure if not _same_hash( snapshot.state.active_config.sha256, expected_active_config_sha256, ) or not _same_hash( snapshot.state.active_secrets.sha256, expected_active_secrets_sha256, ): return None, _candidate_failure('reference', 'schema', 'revision') if required_config: if not snapshot.state.candidate_config.present or not _same_hash( snapshot.state.candidate_config.sha256, expected_candidate_config_sha256, ): return None, _candidate_failure('reference', 'config', 'revision') config_payload = snapshot.candidate_config else: config_payload = snapshot.active_config if required_secrets: if not snapshot.state.candidate_secrets.present or not _same_hash( snapshot.state.candidate_secrets.sha256, expected_candidate_secrets_sha256, ): return None, _candidate_failure('reference', 'secrets', 'revision') secrets_payload = snapshot.candidate_secrets else: secrets_payload = snapshot.active_secrets validated, failure = _validate_payload_pair(config_payload, secrets_payload) if failure is not None: return None, failure return ManagedRuntimeApplyVerification( action=action, state=snapshot.state, config_source='candidate' if required_config else 'active', secrets_source='candidate' if required_secrets else 'active', effective_config_sha256=hashlib.sha256(config_payload).hexdigest(), effective_secrets_sha256=hashlib.sha256(secrets_payload).hexdigest(), ), None except Exception: return None, _candidate_failure() finally: snapshot = None validated = None config_payload = None secrets_payload = None def verify_managed_runtime_candidates( config_path, action, *, expected_active_config_sha256, expected_active_secrets_sha256, expected_candidate_config_sha256=None, expected_candidate_secrets_sha256=None, ): verified, failure = _verify_candidates_inner( config_path, action, expected_active_config_sha256, expected_active_secrets_sha256, expected_candidate_config_sha256, expected_candidate_secrets_sha256, ) config_path = None action = None expected_active_config_sha256 = None expected_active_secrets_sha256 = None expected_candidate_config_sha256 = None expected_candidate_secrets_sha256 = None if failure is not None: raise RuntimeDocumentError( failure[0], failure[1], failure[2], document=failure[3], path=failure[4], ) return verified