import sys import os if __name__ == '__main__': if sys.platform != 'linux' or not os.path.isfile('/.dockerenv') or os.path.abspath(__file__) != '/opt/truf/app/postgres_runtime.py': raise SystemExit('Docker development copy: runtime control is disabled outside the prepared container. See DOCKER_MIGRATION.md.') import runpy runpy.run_path('/opt/truf/app/container_runtime.py')['require_container']() sys.dont_write_bytecode = True if not sys.dont_write_bytecode: raise RuntimeError('PostgreSQL runtime could not disable bytecode writes') import argparse import contextlib import json import re import shlex import signal import socket import stat import subprocess import tempfile import time from concurrent.futures import ThreadPoolExecutor from dataclasses import dataclass from datetime import datetime, timezone from enum import Enum from urllib.parse import quote from db_backend import DatabaseUrlError, canonical_postgres_url, connect_postgres, database_url_from_env, parse_postgres_url from paths import apply_path_config, resolve_postgres_bin_dir, resolve_postgres_data_dir from process_identity import ( ProcessExitedError, ProcessIdentityError, current_process_identity, open_process, ) from runtime_security import ( ClusterAuthorityLock, canonical_cluster_data_directory, canonical_path, harden_private_file, is_reparse_point, private_file_ready, read_private_json, preflight_lifecycle_paths, reject_reparse_components, require_private_directory, require_trusted_native_executable, sha256_file, write_private_json_exclusive, ) CREATE_NO_WINDOW = 0x08000000 IS_WINDOWS = os.name == 'nt' CLUSTER_IDENTITY_SCHEMA = 1 DEFAULT_CONNECT_TIMEOUT_SEC = 5 DEFAULT_QUERY_TIMEOUT_MS = 5000 HOST_AGENT_POSTGRES_SOCKET_DIRECTORY = '/run/truf-postgres' _HOST_AGENT_AUTH_BEGIN = '# BEGIN TRUF HOST AGENT AUTHORITY\n' _HOST_AGENT_AUTH_END = '# END TRUF HOST AGENT AUTHORITY\n' class PostgresState(str, Enum): DISABLED = 'DISABLED' VERIFYING = 'VERIFYING' STARTING = 'STARTING' RECOVERING = 'RECOVERING' STABILIZING = 'STABILIZING' READY = 'READY' BACKOFF = 'BACKOFF' FOREIGN_OR_CONFIG_ERROR = 'FOREIGN_OR_CONFIG_ERROR' STOPPING = 'STOPPING' STOPPED = 'STOPPED' STOP_FAILED = 'STOP_FAILED' def _write_managed_postgres_auth(path, lines): path = os.path.abspath(os.fspath(path)) parent = os.path.dirname(path) reject_reparse_components(parent) details = os.stat(path, follow_symlinks=False) if ( not stat.S_ISREG(details.st_mode) or details.st_nlink != 1 or (not IS_WINDOWS and details.st_uid != os.geteuid()) or stat.S_IMODE(details.st_mode) & 0o077 ): raise ClusterIdentityError('PostgreSQL authentication file is unsafe') descriptor = None try: descriptor = os.open( path, os.O_RDONLY | getattr(os, 'O_CLOEXEC', 0) | getattr(os, 'O_NOFOLLOW', 0), ) before = os.fstat(descriptor) with os.fdopen(descriptor, 'rb') as handle: descriptor = None payload = handle.read(1024 * 1024 + 1) after = os.fstat(handle.fileno()) current = os.stat(path, follow_symlinks=False) identity = lambda item: ( item.st_dev, item.st_ino, item.st_size, getattr(item, 'st_mtime_ns', None), getattr(item, 'st_ctime_ns', None), ) if identity(before) != identity(after) or identity(after) != identity(current): raise OSError('changed') except Exception: raise ClusterIdentityError('PostgreSQL authentication file is unsafe') from None finally: if descriptor is not None: os.close(descriptor) if len(payload) > 1024 * 1024 or b'\x00' in payload: raise ClusterIdentityError('PostgreSQL authentication file is invalid') begin = _HOST_AGENT_AUTH_BEGIN.encode('ascii') end = _HOST_AGENT_AUTH_END.encode('ascii') block = begin + ''.join(line + '\n' for line in lines).encode('ascii') + end if begin in payload or end in payload: if payload.count(begin) != 1 or payload.count(end) != 1 or not payload.startswith(block): raise ClusterIdentityError('PostgreSQL host-agent authority conflicts with managed policy') return stage = path + '.truf-host-agent-stage' descriptor = None try: flags = os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, 'O_NOFOLLOW', 0) descriptor = os.open(stage, flags, 0o600) with os.fdopen(descriptor, 'wb') as handle: descriptor = None handle.write(block) handle.write(payload) handle.flush() os.fsync(handle.fileno()) os.replace(stage, path) directory = os.open(parent, os.O_RDONLY | getattr(os, 'O_DIRECTORY', 0)) try: os.fsync(directory) finally: os.close(directory) except Exception: try: os.unlink(stage) except FileNotFoundError: pass raise ClusterIdentityError('PostgreSQL host-agent authority could not be installed') from None finally: if descriptor is not None: os.close(descriptor) def _configure_host_agent_peer_authority(paths, values): if IS_WINDOWS or values.get('user') != 'truf' or values.get('database') != 'truf': return data_dir = paths['data_dir'] _write_managed_postgres_auth( os.path.join(data_dir, 'pg_hba.conf'), ('local truf truf peer map=truf_host_agent',), ) _write_managed_postgres_auth( os.path.join(data_dir, 'pg_ident.conf'), ('truf_host_agent root truf', 'truf_host_agent truf truf'), ) class ProbeKind(str, Enum): READY = 'READY' RECOVERING = 'RECOVERING' STOPPED = 'STOPPED' OWNED_START_UNCERTAIN = 'OWNED_START_UNCERTAIN' FOREIGN_OR_CONFIG_ERROR = 'FOREIGN_OR_CONFIG_ERROR' @dataclass(frozen=True) class ProbeResult: kind: ProbeKind detail: str = '' postmaster_epoch: str = '' @dataclass(frozen=True) class StartResult: accepted: bool detail: str = '' foreign_or_config_error: bool = False uncertain: bool = False @dataclass(frozen=True) class StopResult: completed: bool stopped: bool detail: str = '' class ClusterIdentityError(ValueError): pass class PostmasterProcessAbsent(ClusterIdentityError): pass class OnlineUnavailable(OSError): pass def utc_now_iso(): return datetime.now(timezone.utc).isoformat(timespec='seconds') def postgres_runtime_paths(config): global_config = (config or {}).get('global') or {} runtime_dir = global_config.get('runtime_dir') if not runtime_dir: root_dir = global_config.get('root_dir') or os.path.dirname(os.path.dirname(os.path.abspath(__file__))) runtime_dir = os.path.join(root_dir, 'runtime') postgres_dir = os.path.join(runtime_dir, 'postgres') data_dir = resolve_postgres_data_dir(global_config, runtime_dir) suffix = '.exe' if os.name == 'nt' else '' bin_dir = resolve_postgres_bin_dir(global_config, runtime_dir) if os.name != 'nt': for path in (postgres_dir, data_dir, bin_dir): reject_reparse_components(path) paths = { 'postgres_dir': canonical_path(postgres_dir), 'bin_dir': canonical_path(bin_dir), 'data_dir': canonical_path(data_dir), 'log_dir': canonical_path(os.path.join(postgres_dir, 'logs')), 'log_path': os.path.join(postgres_dir, 'postgres.log'), 'identity_path': os.path.join(postgres_dir, 'cluster_identity.json'), } for name in ('postgres', 'pg_ctl', 'pg_isready', 'pg_controldata', 'initdb', 'psql'): path = os.path.join(bin_dir, name + suffix) if os.name != 'nt': reject_reparse_components(path) paths[name] = canonical_path(path) return paths def _numbered_log_path(path): base, extension = os.path.splitext(path) sequence = 1 while True: candidate = f'{base}.{sequence:06d}{extension or ".log"}' if not os.path.exists(candidate): return candidate sequence += 1 def rotate_bounded_postgres_startup_log(path, max_bytes, keep): absolute = os.path.abspath(os.fspath(path)) parent = canonical_path(require_private_directory(os.path.dirname(absolute), create=False)) if not os.path.lexists(absolute): return if is_reparse_point(absolute): raise ClusterIdentityError(f'unsafe PostgreSQL startup log: {absolute}') reject_reparse_components(absolute) resolved = canonical_path(absolute) try: contained = os.path.commonpath((parent, resolved)) == parent except ValueError: contained = False if not contained or os.path.dirname(resolved) != parent: raise ClusterIdentityError(f'escaping PostgreSQL startup log: {absolute}') details = os.stat(absolute, follow_symlinks=False) if not stat.S_ISREG(details.st_mode) or details.st_nlink != 1: raise ClusterIdentityError(f'unsafe PostgreSQL startup log: {absolute}') harden_private_file(absolute) verified = os.stat(absolute, follow_symlinks=False) if ( not os.path.samestat(details, verified) or not stat.S_ISREG(verified.st_mode) or verified.st_nlink != 1 or is_reparse_point(absolute) or not private_file_ready(absolute) ): raise ClusterIdentityError(f'PostgreSQL startup log hardening verification failed: {absolute}') if os.path.getsize(absolute) >= max(1, int(max_bytes)): if os.path.getsize(absolute) > max(1, int(max_bytes)): with open(absolute, 'r+b') as handle: handle.truncate(max(1, int(max_bytes))) handle.flush() os.fsync(handle.fileno()) destination = _numbered_log_path(absolute) os.replace(absolute, destination) harden_private_file(destination) base, extension = os.path.splitext(os.path.basename(absolute)) pattern = re.compile(rf'^{re.escape(base)}\.\d{{6}}{re.escape(extension or ".log")}$') rotated = [] with os.scandir(parent) as entries: for entry in entries: if ( pattern.fullmatch(entry.name) and entry.is_file(follow_symlinks=False) and not entry.is_symlink() and not is_reparse_point(entry.path) ): rotated.append(entry.path) rotated.sort(key=lambda value: os.path.getmtime(value), reverse=True) for old in rotated[max(0, int(keep or 0)):]: reject_reparse_components(old) os.remove(old) def prune_postgres_collector_logs(log_dir, keep, max_bytes=0): if not os.path.lexists(log_dir): return 0 verified_log_dir = require_private_directory(log_dir, create=False) verified_root = canonical_path(verified_log_dir) pattern = re.compile(r'postgresql-[0-9]{8}-[0-9]{6}\.log') logs = [] with os.scandir(verified_log_dir) as entries: for entry in entries: if not pattern.fullmatch(entry.name): raise ClusterIdentityError(f'unknown PostgreSQL collector log entry: {entry.path}') if entry.is_symlink() or is_reparse_point(entry.path): raise ClusterIdentityError(f'unsafe PostgreSQL collector log entry: {entry.path}') details = os.stat(entry.path, follow_symlinks=False) if not stat.S_ISREG(details.st_mode) or details.st_nlink != 1: raise ClusterIdentityError(f'unsafe PostgreSQL collector log entry: {entry.path}') if os.path.dirname(canonical_path(entry.path)) != verified_root: raise ClusterIdentityError(f'escaping PostgreSQL collector log entry: {entry.path}') harden_private_file(entry.path) if not private_file_ready(entry.path): raise ClusterIdentityError(f'PostgreSQL collector log is not private: {entry.path}') logs.append(entry.path) logs.sort(key=lambda value: os.path.getmtime(value), reverse=True) removed = 0 keep = max(1, int(keep or 1)) max_bytes = max(0, int(max_bytes or 0)) for index, old in enumerate(logs): size = os.path.getsize(old) if index == 0 and max_bytes and size > max_bytes * 2: raise ClusterIdentityError('active PostgreSQL collector log exceeds the verified rotation bound') if index >= keep or (index > 0 and max_bytes and size > max_bytes): reject_reparse_components(old) os.remove(old) removed += 1 return removed def configured_cluster_values(): return { 'database': str(os.getenv('TRUF_POSTGRES_DB') or 'truf'), 'user': str(os.getenv('TRUF_POSTGRES_USER') or 'truf'), 'port': int(os.getenv('TRUF_POSTGRES_PORT') or 5432), } def canonical_database_url(): values = configured_cluster_values() url = database_url_from_env() if not url and os.getenv('TRUF_POSTGRES_PASSWORD') is not None: password = os.getenv('TRUF_POSTGRES_PASSWORD') or '' url = ( f'postgresql://{quote(values["user"], safe="")}:{quote(password, safe="")}' f'@127.0.0.1:{values["port"]}/{quote(values["database"], safe="")}' ) if not url: return '' return canonical_postgres_url(url, values['database'], values['user'], values['port']) def _postgres_subprocess_environment(overrides=None): environment = os.environ.copy() if os.name != 'nt': environment = {key: value for key, value in environment.items() if not key.upper().startswith('PG')} environment['LC_ALL'] = 'C' environment['LANG'] = 'C' environment.update(overrides or {}) return environment def _run_bounded(command, timeout, capture_output=True, environment_overrides=None): if os.name != 'nt': command = [require_trusted_native_executable(command[0]), *command[1:]] return subprocess.run( command, stdin=subprocess.DEVNULL, stdout=subprocess.PIPE if capture_output else subprocess.DEVNULL, stderr=subprocess.STDOUT if capture_output else subprocess.DEVNULL, text=True, errors='replace', timeout=max(1, float(timeout)), env=_postgres_subprocess_environment(environment_overrides), creationflags=CREATE_NO_WINDOW if os.name == 'nt' else 0, close_fds=True, ) def _system_identifier(paths): result = _run_bounded([paths['pg_controldata'], paths['data_dir']], timeout=15) if result.returncode != 0: raise ClusterIdentityError((result.stdout or '').strip() or 'pg_controldata failed') match = re.search(r'^Database system identifier\s*:\s*(\d+)\s*$', result.stdout or '', re.MULTILINE | re.IGNORECASE) if not match: raise ClusterIdentityError('pg_controldata did not report a system identifier') return match.group(1) def _data_directory_major(paths): version_path = os.path.join(paths['data_dir'], 'PG_VERSION') try: with open(version_path, 'r', encoding='ascii') as handle: cluster_major = handle.read().strip().split('.')[0] except OSError as exc: raise ClusterIdentityError(f'unable to read {version_path}') from exc if not cluster_major.isdigit(): raise ClusterIdentityError(f'invalid PostgreSQL cluster version: {cluster_major!r}') return int(cluster_major) def _bundled_postgres_major(paths): result = _run_bounded([paths['postgres'], '--version'], timeout=10) match = re.search(r'(\d+)(?:\.\d+)?', result.stdout or '') if result.returncode != 0 or not match: raise ClusterIdentityError('unable to determine bundled postgres executable major version') return int(match.group(1)) def _bootstrap_postgres_major(paths): cluster_major = _data_directory_major(paths) if _bundled_postgres_major(paths) != cluster_major: raise ClusterIdentityError('bundled postgres executable does not match the data directory major version') return cluster_major def _listener_present(port, timeout=0.25): try: with socket.create_connection(('127.0.0.1', int(port)), timeout=max(0.05, float(timeout))): return True except OSError: return False def _require_bootstrap_offline(config, paths, values): supervisor = (config or {}).get('supervisor') or {} global_config = (config or {}).get('global') or {} log_dir = supervisor.get('log_dir') or global_config.get('log_dir') or os.path.join(global_config.get('runtime_dir') or '', 'logs') control_dir = supervisor.get('control_dir') or global_config.get('control_dir') or os.path.join(global_config.get('runtime_dir') or '', 'control') instance_path = supervisor.get('instance_file') or os.path.join(control_dir, 'supervisor.instance.json') if os.path.lexists(instance_path): raise ClusterIdentityError(f'refusing bootstrap while supervisor metadata exists: {instance_path}') legacy_instance_path = os.path.join(log_dir, 'supervisor.instance.json') if os.path.abspath(legacy_instance_path) != os.path.abspath(instance_path) and os.path.lexists(legacy_instance_path): raise ClusterIdentityError(f'refusing bootstrap while legacy supervisor metadata exists: {legacy_instance_path}') legacy_pid_path = os.path.join(log_dir, 'supervisor.pid') if os.path.lexists(legacy_pid_path): raise ClusterIdentityError(f'refusing bootstrap while legacy supervisor status metadata exists: {legacy_pid_path}') postmaster_pid = os.path.join(paths['data_dir'], 'postmaster.pid') if os.path.lexists(postmaster_pid): raise ClusterIdentityError(f'refusing bootstrap while postmaster.pid exists: {postmaster_pid}') if _listener_present(values['port']): raise ClusterIdentityError(f'refusing bootstrap while a listener is present on 127.0.0.1:{values["port"]}') def bootstrap_cluster_identity(config): paths = postgres_runtime_paths(config) values = configured_cluster_values() if os.path.lexists(paths['identity_path']): raise ClusterIdentityError('refusing to replace an existing cluster identity') for key in ('data_dir', 'postgres', 'pg_ctl', 'pg_isready', 'pg_controldata'): if key == 'data_dir' and not os.path.isdir(paths[key]): raise ClusterIdentityError(f'bundled PostgreSQL data directory not found: {paths[key]}') if key != 'data_dir' and not os.path.isfile(paths[key]): raise ClusterIdentityError(f'bundled PostgreSQL executable not found: {paths[key]}') if key != 'data_dir' and os.name != 'nt': require_trusted_native_executable(paths[key]) _require_bootstrap_offline(config, paths, values) major = _bootstrap_postgres_major(paths) executable_keys = ('postgres', 'pg_ctl', 'pg_isready', 'pg_controldata') identity = { 'schema': CLUSTER_IDENTITY_SCHEMA, 'private_file_ready': True, 'data_directory': paths['data_dir'], 'pg_major': major, 'executables': {key: paths[key] for key in executable_keys}, 'executable_sha256': {key: sha256_file(paths[key]) for key in executable_keys}, 'database': values['database'], 'user': values['user'], 'port': values['port'], 'system_identifier': _system_identifier(paths), 'created_at': utc_now_iso(), } write_private_json_exclusive(paths['identity_path'], identity) return verify_cluster_identity(config) def _wait_initialization_child(process, *, stop=False): # A direct, unreaped child is the owner, never a PID guessed from PGDATA. # Do not unwind the authority lock or remove its private socket on uncertainty. while True: try: if stop and process.poll() is None: process.send_signal(signal.SIGINT) return process.wait(timeout=60) except (OSError, subprocess.SubprocessError, KeyboardInterrupt, SystemExit): try: print('PostgreSQL initialize-empty FAILED_HOLD: retaining authority until the owned child exits.', flush=True) time.sleep(1) except (OSError, KeyboardInterrupt, SystemExit): pass def initialize_empty(config): """Initialize a prepared, empty, independent POSIX PGDATA and bind it once. The configured user is initdb's bootstrap superuser. No application schema or cutover evidence is installed here; those belong to the ordinary migration. """ if os.name == 'nt' or not hasattr(os, 'geteuid') or os.geteuid() == 0: raise ClusterIdentityError('initialize-empty requires a non-root POSIX runtime user') global_config = (config or {}).get('global') or {} for key in ('runtime_dir', 'postgres_data_dir'): value = global_config.get(key) if not value or not os.path.isabs(value) or any(char in str(value) for char in ('{', '}', '\x00', '\r', '\n')): raise ClusterIdentityError(f'initialize-empty requires an explicit resolved absolute {key}') reject_reparse_components(value) paths = postgres_runtime_paths(config) if paths['data_dir'] != canonical_cluster_data_directory(config): raise ClusterIdentityError('initialize-empty data directory authority mismatch') for key in ('runtime_dir', 'root_dir', 'project_dir'): if global_config.get(key): other = canonical_path(reject_reparse_components(global_config[key])) if os.path.commonpath((other, paths['data_dir'])) in (other, paths['data_dir']): raise ClusterIdentityError('initialize-empty PGDATA must be independent of application and runtime trees') for path in (global_config['runtime_dir'], paths['postgres_dir'], paths['data_dir']): require_private_directory(path, create=False) endpoint_dsn = canonical_database_url() if not endpoint_dsn: raise ClusterIdentityError('initialize-empty requires a canonical managed PostgreSQL DSN') credentials = parse_postgres_url(endpoint_dsn) password = credentials['password'] if not password or any(char in password for char in ('\x00', '\r', '\n')): raise ClusterIdentityError('initialize-empty requires a nonempty single-line PostgreSQL password') values = configured_cluster_values() for key in ('user', 'database'): value = values[key] if not value or len(value.encode('utf-8')) > 63 or any(char in value for char in ('\x00', '\r', '\n')): raise ClusterIdentityError(f'initialize-empty PostgreSQL {key} is invalid') if values['database'] in ('template0', 'template1') or values['user'].startswith('pg_'): raise ClusterIdentityError('initialize-empty cannot use a reserved database or role name') with ClusterAuthorityLock(config, endpoint_dsn=endpoint_dsn): if os.path.lexists(paths['identity_path']): raise ClusterIdentityError('initialize-empty refuses an existing cluster identity') _require_bootstrap_offline(config, paths, values) with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as port_probe: try: port_probe.bind(('127.0.0.1', values['port'])) except OSError as exc: raise ClusterIdentityError('initialize-empty PostgreSQL TCP port is occupied or unavailable') from exc with os.scandir(paths['data_dir']) as entries: if next(entries, None) is not None: raise ClusterIdentityError('initialize-empty refuses nonempty PGDATA; no adoption or repair is performed') for key in ('postgres', 'pg_ctl', 'pg_isready', 'pg_controldata', 'initdb', 'psql'): require_trusted_native_executable(paths[key]) if _bundled_postgres_major(paths) != 16: raise ClusterIdentityError('initialize-empty requires PostgreSQL 16') backend = PostgresBackend(config) backend._prepare_logging() with tempfile.TemporaryDirectory(prefix='init-', dir=paths['postgres_dir']) as private_dir: require_private_directory(private_dir, create=False) password_path = os.path.join(private_dir, 'password') passfile_path = os.path.join(private_dir, 'pgpass') escaped_user = values['user'].replace('\\', '\\\\').replace(':', '\\:') escaped_password = password.replace('\\', '\\\\').replace(':', '\\:') for path, payload in ( (password_path, password + '\n'), (passfile_path, f'*:{values["port"]}:*:{escaped_user}:{escaped_password}\n'), ): descriptor = os.open(path, os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, 'O_NOFOLLOW', 0), 0o600) with os.fdopen(descriptor, 'w', encoding='utf-8', newline='\n') as handle: handle.write(payload) handle.flush() os.fsync(handle.fileno()) if not private_file_ready(path): raise ClusterIdentityError('initialize-empty password file is not private') with open(paths['log_path'], 'ab', buffering=0) as startup_log: initdb = subprocess.Popen( [require_trusted_native_executable(paths['initdb']), '-D', paths['data_dir'], '--encoding=UTF8', '--locale=C', '--auth-local=scram-sha-256', '--auth-host=scram-sha-256', '--username=' + values['user'], '--pwfile=' + password_path], stdin=subprocess.DEVNULL, stdout=startup_log, stderr=subprocess.STDOUT, close_fds=True, start_new_session=True, env=_postgres_subprocess_environment(), ) try: initdb_result = initdb.wait(timeout=120) finally: _wait_initialization_child(initdb) if initdb_result != 0: raise ClusterIdentityError('initdb failed; PGDATA is retained for review, see the private PostgreSQL log') _configure_host_agent_peer_authority(paths, values) system_identifier = _system_identifier(paths) socket_setting = 'unix_socket_directories="' + private_dir.replace('"', '""') + '"' backend._start_requested_wall_time = time.time() process = subprocess.Popen( [require_trusted_native_executable(paths['postgres']), '-D', paths['data_dir'], '-c', 'data_directory=' + paths['data_dir'], '-c', 'listen_addresses=', '-c', socket_setting, '-c', 'unix_socket_permissions=0700', '-c', f'port={values["port"]}', '-c', 'logging_collector=off'], stdin=subprocess.DEVNULL, stdout=startup_log, stderr=subprocess.STDOUT, close_fds=True, start_new_session=True, env=_postgres_subprocess_environment(), ) try: deadline = time.monotonic() + 30 while True: if process.poll() is not None: raise ClusterIdentityError('temporary PostgreSQL exited before initialization') ready = _run_bounded( [paths['pg_isready'], '-h', private_dir, '-p', str(values['port']), '-U', values['user'], '-d', 'postgres', '-t', '1'], timeout=5, ) if ready.returncode == 0: break if time.monotonic() >= deadline: raise ClusterIdentityError('temporary PostgreSQL socket readiness timed out') time.sleep(0.1) retained = backend._open_postmaster({ 'data_directory': paths['data_dir'], 'executables': {'postgres': paths['postgres']}, }, allow_new_after_start=True) try: if retained.pid != process.pid: raise ClusterIdentityError('temporary postmaster is not the directly owned child') finally: retained.close() def execute(database, statement): result = _run_bounded( [paths['psql'], '-X', '-w', '-A', '-t', '--set=ON_ERROR_STOP=1', '-h', private_dir, '-p', str(values['port']), '-U', values['user'], '-d', database, '-c', statement], timeout=30, environment_overrides={'PGPASSFILE': passfile_path, 'PGCONNECT_TIMEOUT': '5'}, ) if result.returncode != 0: raise ClusterIdentityError('temporary PostgreSQL SQL setup failed: ' + (result.stdout or '').strip()) return result.stdout or '' online = json.loads(execute('postgres', """ SELECT pg_catalog.json_build_object( 'data_directory', pg_catalog.current_setting('data_directory'), 'system_identifier', (SELECT system_identifier::text FROM pg_catalog.pg_control_system()), 'user', CURRENT_USER, 'port', pg_catalog.current_setting('port')::integer, 'listen_addresses', pg_catalog.current_setting('listen_addresses'), 'unix_socket', pg_catalog.inet_server_addr() IS NULL) """)) if online != { 'data_directory': paths['data_dir'], 'system_identifier': system_identifier, 'user': values['user'], 'port': values['port'], 'listen_addresses': '', 'unix_socket': True, }: raise ClusterIdentityError('temporary PostgreSQL online identity mismatch') role = '"' + values['user'].replace('"', '""') + '"' database = '"' + values['database'].replace('"', '""') + '"' if values['database'] != 'postgres': execute('postgres', f"CREATE DATABASE {database} OWNER {role} ENCODING 'UTF8' TEMPLATE template0") execute(values['database'], f'REVOKE CREATE ON SCHEMA public FROM PUBLIC; ALTER ROLE {role} SET search_path TO public') finally: stop_result = _wait_initialization_child(process, stop=True) if stop_result != 0: raise ClusterIdentityError('temporary PostgreSQL did not exit cleanly; no cluster identity was bound') _require_bootstrap_offline(config, paths, values) control = _run_bounded([paths['pg_controldata'], paths['data_dir']], timeout=15) if control.returncode != 0 or not re.search(r'^Database cluster state\s*:\s*shut down\s*$', control.stdout or '', re.MULTILINE): raise ClusterIdentityError('initialize-empty could not confirm a cleanly stopped cluster') return bootstrap_cluster_identity(config) def _validated_identity(value): if not isinstance(value, dict) or value.get('schema') != CLUSTER_IDENTITY_SCHEMA: raise ClusterIdentityError('unsupported cluster identity schema') if value.get('private_file_ready') is not True: raise ClusterIdentityError('cluster identity is not private-file-ready') required_strings = ('data_directory', 'database', 'user', 'system_identifier', 'created_at') for key in required_strings: if not isinstance(value.get(key), str) or not value[key]: raise ClusterIdentityError(f'invalid cluster identity field: {key}') if not value['system_identifier'].isdigit(): raise ClusterIdentityError('invalid cluster system identifier') try: major = int(value.get('pg_major')) port = int(value.get('port')) except (TypeError, ValueError) as exc: raise ClusterIdentityError('invalid cluster version or port') from exc executables = value.get('executables') hashes = value.get('executable_sha256') keys = ('postgres', 'pg_ctl', 'pg_isready', 'pg_controldata') if not isinstance(executables, dict) or not isinstance(hashes, dict): raise ClusterIdentityError('invalid cluster executable identity') normalized = dict(value) normalized['pg_major'] = major normalized['port'] = port normalized['data_directory'] = canonical_path(value['data_directory']) normalized['executables'] = {key: canonical_path(executables.get(key, '')) for key in keys} normalized['executable_sha256'] = {key: str(hashes.get(key) or '') for key in keys} return normalized def verify_cluster_identity(config, *, verify_offline_system_identifier=True): paths = postgres_runtime_paths(config) try: identity = _validated_identity(read_private_json(paths['identity_path'])) except OSError as exc: raise ClusterIdentityError(str(exc)) from exc values = configured_cluster_values() expected = { 'data_directory': paths['data_dir'], 'database': values['database'], 'user': values['user'], 'port': values['port'], } for key, expected_value in expected.items(): if identity[key] != expected_value: raise ClusterIdentityError(f'cluster identity {key} mismatch') for key in ('postgres', 'pg_ctl', 'pg_isready', 'pg_controldata'): if identity['executables'][key] != paths[key] or not os.path.isfile(paths[key]): raise ClusterIdentityError(f'cluster executable path mismatch: {key}') if os.name != 'nt': require_trusted_native_executable(paths[key]) if identity['executable_sha256'][key] != sha256_file(paths[key]): raise ClusterIdentityError(f'cluster executable hash mismatch: {key}') if identity['pg_major'] != _data_directory_major(paths): raise ClusterIdentityError('cluster PostgreSQL major version changed') if verify_offline_system_identifier and identity['system_identifier'] != _system_identifier(paths): raise ClusterIdentityError('offline cluster system identifier changed') return identity def _database_url_matches_identity(url, identity): try: canonical_postgres_url(url, identity['database'], identity['user'], identity['port']) except (TypeError, ValueError): return False return True @dataclass(frozen=True) class PostmasterPidRecord: pid: int data_directory: str start_time: float def _parse_postmaster_pid_record(data_dir): path = os.path.join(data_dir, 'postmaster.pid') try: with open(path, 'r', encoding='ascii') as handle: lines = [handle.readline().strip() for _ in range(3)] pid = int(lines[0]) start_time = float(lines[2]) if pid <= 0 or not lines[1]: return None return PostmasterPidRecord(pid, canonical_path(lines[1]), start_time) except (OSError, ValueError, IndexError): return None def _parse_postmaster_pid(data_dir): record = _parse_postmaster_pid_record(data_dir) return record.pid if record else None def _timestamp(value): if hasattr(value, 'timestamp'): return float(value.timestamp()) return datetime.fromisoformat(str(value).replace('Z', '+00:00')).timestamp() class PostgresBackend: """Bounded ordinary-subprocess operations used only by one controller worker.""" def __init__( self, config, connect_timeout_sec=DEFAULT_CONNECT_TIMEOUT_SEC, query_timeout_ms=DEFAULT_QUERY_TIMEOUT_MS, stop_timeout_sec=60, start_settle_timeout_sec=30, log_max_mb=64, log_keep=24, ): self.config = config self.paths = postgres_runtime_paths(config) self.connect_timeout_sec = max(1, int(connect_timeout_sec)) self.query_timeout_ms = max(1000, int(query_timeout_ms)) self.stop_timeout_sec = max(5, int(stop_timeout_sec)) self.start_settle_timeout_sec = min(120.0, max(0.1, float(start_settle_timeout_sec))) self.log_max_mb = max(1, int(log_max_mb)) self.log_keep = max(1, int(log_keep)) self._expected_process = None self._start_requested_wall_time = None self._accepted_start_at_monotonic = None self._started_postmaster_observed = False self._verified_identity = None self._authenticated_postmaster = None @property def owns_start(self): return ( getattr(self, '_start_requested_wall_time', None) is not None or getattr(self, '_accepted_start_at_monotonic', None) is not None or getattr(self, '_started_postmaster_observed', False) ) def _prepare_logging(self): log_dir = self.paths.get('log_dir') if not log_dir: return require_private_directory(os.path.dirname(os.path.abspath(self.paths['log_path'])), create=False) require_private_directory(log_dir, create=True) rotate_bounded_postgres_startup_log( self.paths['log_path'], max(1, int(getattr(self, 'log_max_mb', 64))) * 1024 * 1024, max(1, int(getattr(self, 'log_keep', 24))), ) if not os.path.lexists(self.paths['log_path']): with open(self.paths['log_path'], 'xb'): pass harden_private_file(self.paths['log_path']) prune_postgres_collector_logs( log_dir, max(1, int(getattr(self, 'log_keep', 24))), max(1, int(getattr(self, 'log_max_mb', 64))) * 1024 * 1024, ) def _prune_logging(self): log_dir = self.paths.get('log_dir') if log_dir and os.path.lexists(log_dir): prune_postgres_collector_logs( log_dir, max(1, int(getattr(self, 'log_keep', 24))), max(1, int(getattr(self, 'log_max_mb', 64))) * 1024 * 1024, ) def _root_job_error(self): if os.name != 'nt': return None try: identity = current_process_identity() except (OSError, ValueError) as exc: return f'unable to verify root supervisor Job membership: {exc}' if identity.in_job is not False: return 'root supervisor is inside an application or unknown Windows Job' return None def _remember_process(self, process): current = self._expected_process if (current and current.is_running() and current.pid == process.pid and current.identity.creation_time == process.identity.creation_time): process.close() return current if current: current.close() self._expected_process = process return process def _forget_dead_process(self): if self._expected_process and not self._expected_process.is_running(): self._expected_process.close() self._expected_process = None self._authenticated_postmaster = None def _expected_is_live(self, identity): self._forget_dead_process() process = self._expected_process if not process or not process.is_running(): return False pid = _parse_postmaster_pid(identity['data_directory']) return pid == process.pid def _open_postmaster(self, identity, allow_new_after_start=False): record = _parse_postmaster_pid_record(identity['data_directory']) if not record: raise ClusterIdentityError('postmaster.pid is absent or invalid') if record.data_directory != identity['data_directory']: raise ClusterIdentityError('postmaster.pid data directory does not match bound cluster') try: process = open_process(record.pid) except ProcessExitedError as exc: raise PostmasterProcessAbsent( f'postmaster.pid references exited process {record.pid}' ) from exc except ProcessIdentityError as exc: cause = exc.__cause__ winerror = getattr(cause, 'winerror', None) or getattr(exc, 'winerror', None) errno_value = getattr(cause, 'errno', None) or getattr(exc, 'errno', None) if (IS_WINDOWS and winerror in (87, 1168)) or ( not IS_WINDOWS and errno_value in (2, 3) ): raise PostmasterProcessAbsent( f'postmaster.pid references absent process {record.pid}' ) from exc raise ClusterIdentityError(str(exc)) from exc if process.identity.executable != identity['executables']['postgres']: process.close() raise ClusterIdentityError('postmaster executable path mismatch') if os.name == 'nt' and process.identity.in_job is not False: process.close() raise ClusterIdentityError('verified postmaster is inside an application or unknown Windows Job') if IS_WINDOWS and abs(process.identity.creation_time_unix - record.start_time) > 10.0: process.close() raise ClusterIdentityError('postmaster.pid start time does not match the live process') if IS_WINDOWS and allow_new_after_start: threshold = float(self._start_requested_wall_time or 0) - 5.0 if not threshold or process.identity.creation_time_unix < threshold: process.close() raise ClusterIdentityError('postmaster creation time does not match this start request') return process def _online_query(self, identity): url = database_url_from_env() try: url = canonical_postgres_url(url, identity['database'], identity['user'], identity['port']) except DatabaseUrlError as exc: raise ClusterIdentityError(str(exc)) from exc connection = None try: connection = connect_postgres( url, connect_timeout_sec=self.connect_timeout_sec, statement_timeout_ms=self.query_timeout_ms, lock_timeout_ms=min(2000, self.query_timeout_ms), idle_in_transaction_timeout_ms=self.query_timeout_ms, tcp_user_timeout_ms=self.query_timeout_ms, ) row = connection.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, pg_catalog.pg_is_in_recovery() AS in_recovery, pg_catalog.pg_postmaster_start_time() AS postmaster_start_time""" ).fetchone() system_identifier = None control_system_available = False if identity['pg_major'] >= 10: availability = connection.execute( """SELECT pg_catalog.to_regprocedure('pg_catalog.pg_control_system()') IS NOT NULL AS present, CASE WHEN pg_catalog.to_regprocedure('pg_catalog.pg_control_system()') IS NULL THEN false ELSE pg_catalog.has_function_privilege( CURRENT_USER, pg_catalog.to_regprocedure('pg_catalog.pg_control_system()'), 'EXECUTE' ) END AS permitted""" ).fetchone() control_system_available = bool(availability and availability['present'] and availability['permitted']) if control_system_available: system_row = connection.execute('SELECT system_identifier::text AS system_identifier FROM pg_catalog.pg_control_system()').fetchone() system_identifier = system_row['system_identifier'] if system_row else None except ClusterIdentityError: raise except Exception as exc: raise OnlineUnavailable(str(exc)) from exc finally: if connection is not None: try: connection.close() except Exception: pass if not row: raise ClusterIdentityError('authenticated readiness query returned no row') checks = { 'database': (str(row['database']), identity['database']), 'user': (str(row['user_name']), identity['user']), 'data_directory': (canonical_path(row['data_directory']), identity['data_directory']), 'port': (int(row['port']), identity['port']), } for name, (actual, expected) in checks.items(): if actual != expected: raise ClusterIdentityError(f'online cluster {name} mismatch') if control_system_available and str(system_identifier or '') != identity['system_identifier']: raise ClusterIdentityError('online cluster system identifier mismatch') process = self._open_postmaster(identity) try: if IS_WINDOWS and abs(process.identity.creation_time_unix - _timestamp(row['postmaster_start_time'])) > 10.0: raise ClusterIdentityError('postmaster process creation time mismatch') self._remember_process(process) # A CLI handoff may lose SQL readiness without losing this exact process. if control_system_available: self._authenticated_postmaster = (dict(identity), self._expected_process) except BaseException: if process is not self._expected_process: process.close() raise return bool(row['in_recovery']), self._expected_process.identity.creation_time def _hint_running(self, identity): try: result = _run_bounded([identity['executables']['pg_ctl'], '-D', identity['data_directory'], 'status'], timeout=10) return result.returncode == 0 except (OSError, subprocess.SubprocessError): return False def probe(self): root_error = self._root_job_error() if root_error: return ProbeResult(ProbeKind.FOREIGN_OR_CONFIG_ERROR, root_error) try: identity = verify_cluster_identity(self.config, verify_offline_system_identifier=False) except (OSError, ValueError, subprocess.SubprocessError) as exc: self._verified_identity = None return ProbeResult(ProbeKind.FOREIGN_OR_CONFIG_ERROR, str(exc)) self._verified_identity = identity try: self._prune_logging() except (OSError, ValueError) as exc: return ProbeResult(ProbeKind.FOREIGN_OR_CONFIG_ERROR, f'PostgreSQL log retention validation failed: {exc}') try: in_recovery, postmaster_epoch = self._online_query(identity) if getattr(self, '_accepted_start_at_monotonic', None) is not None: self._accepted_start_at_monotonic = None self._started_postmaster_observed = True if in_recovery: return ProbeResult( ProbeKind.RECOVERING, 'authenticated expected PostgreSQL reports recovery mode', postmaster_epoch, ) return ProbeResult(ProbeKind.READY, 'authenticated cluster identity verified', postmaster_epoch) except PostmasterProcessAbsent as exc: unavailable_detail = str(exc) except ClusterIdentityError as exc: return ProbeResult(ProbeKind.FOREIGN_OR_CONFIG_ERROR, str(exc)) except OnlineUnavailable as exc: unavailable_detail = str(exc) if self._expected_is_live(identity): if getattr(self, '_accepted_start_at_monotonic', None) is not None: self._accepted_start_at_monotonic = None self._started_postmaster_observed = True return ProbeResult( ProbeKind.RECOVERING, f'expected postmaster is live but unavailable: {unavailable_detail}', self._expected_process.identity.creation_time, ) pid = _parse_postmaster_pid(identity['data_directory']) if pid and self._start_requested_wall_time: try: self._remember_process(self._open_postmaster(identity, allow_new_after_start=True)) self._accepted_start_at_monotonic = None self._started_postmaster_observed = True return ProbeResult(ProbeKind.RECOVERING, f'new expected postmaster is still starting: {unavailable_detail}') except PostmasterProcessAbsent: pid = None except ClusterIdentityError as exc: return ProbeResult(ProbeKind.FOREIGN_OR_CONFIG_ERROR, str(exc)) if pid: try: self._remember_process(self._open_postmaster(identity)) return ProbeResult( ProbeKind.RECOVERING, f'bound bundled postmaster is live and awaiting authenticated identity: {unavailable_detail}', self._expected_process.identity.creation_time, ) except PostmasterProcessAbsent: pid = None except ClusterIdentityError as exc: return ProbeResult(ProbeKind.FOREIGN_OR_CONFIG_ERROR, str(exc)) if _listener_present(identity['port']): return ProbeResult(ProbeKind.FOREIGN_OR_CONFIG_ERROR, f'unauthenticated or foreign listener on 127.0.0.1:{identity["port"]}') if self._hint_running(identity): return ProbeResult(ProbeKind.FOREIGN_OR_CONFIG_ERROR, 'pg_ctl reports an unverified running postmaster') if getattr(self, '_accepted_start_at_monotonic', None) is not None: elapsed = time.monotonic() - self._accepted_start_at_monotonic if elapsed >= float(getattr(self, 'start_settle_timeout_sec', 30.0)): return ProbeResult( ProbeKind.OWNED_START_UNCERTAIN, 'controller-owned PostgreSQL start did not materialize within the bounded settle interval', ) return ProbeResult( ProbeKind.RECOVERING, 'accepted PostgreSQL start remains unobservable; start ownership is retained', ) try: verify_cluster_identity(self.config) except (OSError, ValueError, subprocess.SubprocessError) as exc: self._verified_identity = None return ProbeResult(ProbeKind.FOREIGN_OR_CONFIG_ERROR, str(exc)) return ProbeResult(ProbeKind.STOPPED, 'bound cluster is offline') def start(self): probe = self.probe() if probe.kind == ProbeKind.FOREIGN_OR_CONFIG_ERROR: return StartResult(False, probe.detail, foreign_or_config_error=True) if probe.kind != ProbeKind.STOPPED: return StartResult(False, f'start refused while cluster state is {probe.kind.value}') identity = self._verified_identity if not identity: return StartResult(False, 'verified cluster identity was not retained', foreign_or_config_error=True) self._start_requested_wall_time = time.time() self._accepted_start_at_monotonic = None self._started_postmaster_observed = False try: self._prepare_logging() except (OSError, ValueError) as exc: self._start_requested_wall_time = None return StartResult(False, f'PostgreSQL logging setup failed: {exc}', foreign_or_config_error=True) port = str(identity['port']) host_agent_authority = ( not IS_WINDOWS and identity.get('user') == 'truf' and identity.get('database') == 'truf' ) if host_agent_authority: require_private_directory( HOST_AGENT_POSTGRES_SOCKET_DIRECTORY, create=False, ) _configure_host_agent_peer_authority(self.paths, identity) options = [ '-c', f'data_directory={identity["data_directory"]}', '-c', 'listen_addresses=127.0.0.1', '-c', f'port={port}', ] if not IS_WINDOWS: options.extend([ '-c', 'unix_socket_directories=' + ( HOST_AGENT_POSTGRES_SOCKET_DIRECTORY if host_agent_authority else '' ), '-c', 'unix_socket_permissions=0700', ]) if self.paths.get('log_dir'): options.extend([ '-c', 'logging_collector=on', '-c', 'log_destination=stderr', '-c', f'log_directory={self.paths["log_dir"]}', '-c', 'log_filename=postgresql-%Y%m%d-%H%M%S.log', '-c', 'log_rotation_age=60', '-c', f'log_rotation_size={max(1, int(getattr(self, "log_max_mb", 64)))}MB', '-c', 'log_truncate_on_rotation=on', '-c', 'log_file_mode=0600', ]) # Retain launch ownership even when an interrupt prevents a launcher # result from reaching us. Only stop() may settle an uncertain launch. self._accepted_start_at_monotonic = time.monotonic() if IS_WINDOWS: try: with open(self.paths['log_path'], 'ab', buffering=0) as startup_log: process = subprocess.Popen( [identity['executables']['postgres'], '-D', identity['data_directory'], *options], stdin=subprocess.DEVNULL, stdout=startup_log, stderr=subprocess.STDOUT, close_fds=True, creationflags=CREATE_NO_WINDOW, ) time.sleep(0.1) exit_code = process.poll() except (OSError, subprocess.SubprocessError) as exc: return StartResult(True, f'detached postgres launch outcome is uncertain: {exc}', uncertain=True) if exit_code is not None: self._start_requested_wall_time = None self._accepted_start_at_monotonic = None return StartResult(False, f'detached postgres exited during launch with code {exit_code}') self._accepted_start_at_monotonic = time.monotonic() return StartResult(True, 'detached postgres start request accepted') option_text = shlex.join(options) command = [ identity['executables']['pg_ctl'], '-D', identity['data_directory'], '-l', self.paths['log_path'], '-o', option_text, 'start', '-W', ] try: result = _run_bounded(command, timeout=30, capture_output=False) except subprocess.TimeoutExpired as exc: self._accepted_start_at_monotonic = time.monotonic() return StartResult( True, f'pg_ctl start timed out after {exc.timeout}s; start outcome is uncertain', uncertain=True, ) except (OSError, subprocess.SubprocessError) as exc: return StartResult(True, f'pg_ctl launch outcome is uncertain: {exc}', uncertain=True) if result.returncode != 0: return StartResult(True, (result.stdout or '').strip() or 'pg_ctl start outcome is uncertain', uncertain=True) self._accepted_start_at_monotonic = time.monotonic() return StartResult(True, 'pg_ctl accepted the start request') def _stop_retained_postmaster(self, identity): process = self._expected_process pid = _parse_postmaster_pid(identity['data_directory']) if not process or not process.is_running() or process.pid != pid: return StopResult(False, False, 'retained postmaster identity no longer matches postmaster.pid') command = [ identity['executables']['pg_ctl'], '-D', identity['data_directory'], 'stop', '-m', 'fast', '-w', '-t', str(self.stop_timeout_sec), ] try: result = _run_bounded(command, timeout=self.stop_timeout_sec + 10) except (OSError, subprocess.SubprocessError) as exc: return StopResult(False, False, f'bounded PostgreSQL stop failed: {exc}') if result.returncode != 0: return StopResult(False, False, (result.stdout or '').strip() or 'pg_ctl stop failed') if process.is_running(): return StopResult(False, False, 'pg_ctl returned success but the retained postmaster is still running') process.close() self._expected_process = None self._authenticated_postmaster = None self._accepted_start_at_monotonic = None self._started_postmaster_observed = False self._start_requested_wall_time = None return StopResult(True, True, 'identity-verified PostgreSQL stop completed') def stop(self): probe = self.probe() accepted_at = getattr(self, '_accepted_start_at_monotonic', None) if accepted_at is not None and not self._expected_process: deadline = time.monotonic() + float(getattr(self, 'start_settle_timeout_sec', 30.0)) while not self._expected_process: if probe.kind == ProbeKind.FOREIGN_OR_CONFIG_ERROR: return StopResult(False, False, f'uncertain PostgreSQL start could not be identity-stopped: {probe.detail}') remaining = deadline - time.monotonic() if remaining <= 0: break time.sleep(min(0.1, remaining)) probe = self.probe() if not self._expected_process: identity = self._verified_identity if ( identity and probe.kind in (ProbeKind.STOPPED, ProbeKind.OWNED_START_UNCERTAIN) and not _listener_present(identity['port']) and not self._hint_running(identity) ): self._accepted_start_at_monotonic = None self._started_postmaster_observed = False self._start_requested_wall_time = None return StopResult( True, True, 'accepted PostgreSQL start did not materialize during bounded compensation probes', ) return StopResult( False, False, 'accepted PostgreSQL start remained unobservable after bounded shutdown probes; stopped state is uncertain', ) if probe.kind == ProbeKind.FOREIGN_OR_CONFIG_ERROR: return StopResult(False, False, f'identity-verified stop refused: {probe.detail}') if accepted_at is not None and self._expected_process: identity = self._verified_identity if not identity: return StopResult(False, False, 'uncertain PostgreSQL start retained a process without verified cluster identity') return self._stop_retained_postmaster(identity) if probe.kind == ProbeKind.STOPPED: return StopResult(True, True, 'cluster already stopped') if probe.kind == ProbeKind.RECOVERING: identity = self._verified_identity if ( (accepted_at is not None or getattr(self, '_started_postmaster_observed', False) or getattr(self, '_authenticated_postmaster', None) == (identity, self._expected_process)) and identity and self._expected_process ): return self._stop_retained_postmaster(identity) return StopResult(False, False, 'live expected recovery was left running because online identity is unavailable') if probe.kind != ProbeKind.READY: return StopResult(False, False, f'identity-verified stop refused: {probe.detail}') identity = self._verified_identity if not identity: return StopResult(False, False, 'verified cluster identity was not retained') return self._stop_retained_postmaster(identity) def close(self): if self._expected_process: self._expected_process.close() self._expected_process = None self._authenticated_postmaster = None class PostgresController: """Nonblocking PostgreSQL state machine. tick() only polls/submits one worker operation.""" def __init__( self, backend, enabled=True, health_interval_sec=5, stable_ready_interval_sec=60, ready_loss_grace_sec=0, backoff_base_sec=30, backoff_max_sec=600, executor=None, authority_check=None, shutdown_timeout_sec=120, ): self.backend = backend self.state = PostgresState.VERIFYING if enabled else PostgresState.DISABLED self.health_interval_sec = max(0.1, float(health_interval_sec)) self.stable_ready_interval_sec = max(0.0, float(stable_ready_interval_sec)) self.ready_loss_grace_sec = max(0.0, float(ready_loss_grace_sec)) self.backoff_base_sec = max(30.0, float(backoff_base_sec)) self.backoff_max_sec = min(600.0, max(self.backoff_base_sec, float(backoff_max_sec))) self.failures = 0 self.detail = '' self.next_action_at = 0.0 self.stable_since = None self._stable_epoch = '' self._ready_epoch = '' self._ready_loss_since = None self._executor = executor or ThreadPoolExecutor(max_workers=1, thread_name_prefix='postgres-runtime') self._owns_executor = executor is None self._future = None self._operation = None self._operation_generation = None self._generation = 0 self._shutdown_requested = False self._stop_submitted = False self._lifecycle_inert = False self._automatic_inhibited = False self._owned_start = False self._compensating_start_failure = False self._authority_release_safe = True self.shutdown_timeout_sec = max(0.1, float(shutdown_timeout_sec)) self.authority_check = authority_check @property def ready(self): return self.state == PostgresState.READY @property def terminal(self): return self.state in (PostgresState.DISABLED, PostgresState.STOPPED, PostgresState.STOP_FAILED) @property def stop_succeeded(self): return self.state in (PostgresState.DISABLED, PostgresState.STOPPED) @property def authority_release_safe(self): return self._authority_release_safe and self._future is None and self.stop_succeeded @property def has_inflight_start(self): return self._operation == 'start' and self._future is not None @property def lifecycle_action_required(self): """Whether shutdown must preserve or complete controller-owned lifecycle work.""" if self.state in (PostgresState.DISABLED, PostgresState.STOPPED): return False return ( self._owned_start or self.has_inflight_start or self._compensating_start_failure or self._operation == 'stop' or self._shutdown_requested ) def _submit(self, operation): if self._future is not None: return False if operation == 'start' and self.authority_check is not None and not self.authority_check(): self._automatic_inhibited = True self._lifecycle_inert = True self.state = PostgresState.FOREIGN_OR_CONFIG_ERROR self.detail = 'runtime authority drifted before PostgreSQL start' return False self._generation += 1 generation = self._generation function = getattr(self.backend, operation) self._future = self._executor.submit(function) self._operation = operation self._operation_generation = generation if operation == 'stop': self._stop_submitted = True return True def _enter_backoff(self, now, detail): self.failures += 1 delay = min(self.backoff_base_sec * (2 ** min(self.failures - 1, 20)), self.backoff_max_sec) self.state = PostgresState.BACKOFF self.detail = detail self.next_action_at = now + delay self.stable_since = None self._stable_epoch = '' self._ready_epoch = '' self._owned_start = False def _handle_probe(self, result, now): if not isinstance(result, ProbeResult): self.state = PostgresState.FOREIGN_OR_CONFIG_ERROR self.detail = 'invalid PostgreSQL probe result' self._lifecycle_inert = True return if result.kind == ProbeKind.FOREIGN_OR_CONFIG_ERROR: self.state = PostgresState.FOREIGN_OR_CONFIG_ERROR self.detail = result.detail self._lifecycle_inert = True self.stable_since = None self._stable_epoch = '' self._ready_epoch = '' self._ready_loss_since = None self.next_action_at = now + self.health_interval_sec return if result.kind == ProbeKind.OWNED_START_UNCERTAIN: if not self._owned_start: self.state = PostgresState.FOREIGN_OR_CONFIG_ERROR self.detail = 'unowned PostgreSQL start uncertainty was reported' self._lifecycle_inert = True return self._compensating_start_failure = True self._authority_release_safe = False self.state = PostgresState.STOPPING self.detail = result.detail or 'controller-owned start requires bounded compensation' self._stop_submitted = False if not self._submit('stop'): self.state = PostgresState.STOP_FAILED self.detail = 'unable to submit identity-safe compensation for an uncertain owned start' return if result.kind == ProbeKind.RECOVERING: if ( self.state == PostgresState.READY and self.ready_loss_grace_sec > 0 and result.postmaster_epoch and result.postmaster_epoch == self._ready_epoch ): if self._ready_loss_since is None: self._ready_loss_since = now if now - self._ready_loss_since < self.ready_loss_grace_sec: self.detail = f'transient readiness loss: {result.detail}' self.next_action_at = now + self.health_interval_sec return self._ready_loss_since = None self.state = PostgresState.RECOVERING self.detail = result.detail # A controller-owned start may pass through recovery without becoming # a foreign adoption. A process found before our own start remains # observational and lifecycle-inert. self._lifecycle_inert = not self._owned_start self.stable_since = None self._stable_epoch = '' self._ready_epoch = '' self.next_action_at = now + self.health_interval_sec return if result.kind == ProbeKind.STOPPED: self._ready_loss_since = None if self._lifecycle_inert and not self._owned_start: self.state = PostgresState.FOREIGN_OR_CONFIG_ERROR self.detail = (result.detail or 'cluster is offline') + '; automatic lifecycle action remains inhibited' self.next_action_at = now + self.health_interval_sec return if self.state in (PostgresState.STARTING, PostgresState.STABILIZING, PostgresState.READY, PostgresState.RECOVERING): self._enter_backoff(now, result.detail or 'PostgreSQL died before stable readiness') else: self.state = PostgresState.VERIFYING self.detail = result.detail self._submit('start') return # Readiness verifies an observed cluster but does not adopt its lifecycle. self._ready_loss_since = None self._lifecycle_inert = not self._owned_start if self.state == PostgresState.READY and ( not result.postmaster_epoch or result.postmaster_epoch == self._ready_epoch ): self.detail = result.detail self.next_action_at = now + self.health_interval_sec return if ( self.state != PostgresState.STABILIZING or self.stable_since is None or (result.postmaster_epoch and result.postmaster_epoch != self._stable_epoch) ): self.state = PostgresState.STABILIZING self.stable_since = now self._stable_epoch = result.postmaster_epoch self.detail = result.detail if self.stable_ready_interval_sec == 0: self.state = PostgresState.READY self.failures = 0 self._ready_epoch = result.postmaster_epoch elif now - self.stable_since >= self.stable_ready_interval_sec: self.state = PostgresState.READY self.failures = 0 self._ready_epoch = result.postmaster_epoch self.detail = result.detail self.next_action_at = now + self.health_interval_sec def _handle_start(self, result, now): if not isinstance(result, StartResult): self._enter_backoff(now, 'invalid PostgreSQL start result') return if result.foreign_or_config_error: self.state = PostgresState.FOREIGN_OR_CONFIG_ERROR self.detail = result.detail self._lifecycle_inert = True self.stable_since = None self._stable_epoch = '' self._ready_epoch = '' self.next_action_at = now + self.health_interval_sec return if not result.accepted: self._enter_backoff(now, result.detail or 'PostgreSQL start request failed') return self.state = PostgresState.STARTING self._owned_start = True self._lifecycle_inert = False self.detail = result.detail self.next_action_at = now def _consume_completion(self, now): if self._future is None or not self._future.done(): return False future = self._future operation = self._operation generation = self._operation_generation self._future = None self._operation = None self._operation_generation = None try: result = future.result() except BaseException as exc: if generation != self._generation: if self._shutdown_requested: self.state = PostgresState.STOPPING return True if self._shutdown_requested or operation == 'stop': self.state = PostgresState.STOP_FAILED self.detail = f'bounded PostgreSQL stop operation failed: {exc}' elif operation == 'probe': self.state = PostgresState.FOREIGN_OR_CONFIG_ERROR self.detail = f'PostgreSQL verification failed: {exc}' self._lifecycle_inert = True else: self._enter_backoff(now, f'PostgreSQL start operation failed: {exc}') return True if generation != self._generation: if operation == 'start' and isinstance(result, StartResult): if result.accepted: self._shutdown_requested = True self._automatic_inhibited = True self.state = PostgresState.STOPPING self._owned_start = True self.detail = 'late accepted PostgreSQL start requires identity-safe compensating stop' if not self._submit('stop'): self.state = PostgresState.STOP_FAILED self.detail = 'unable to schedule identity-safe stop for late accepted PostgreSQL start' elif self._shutdown_requested: self._automatic_inhibited = True self._lifecycle_inert = True self._owned_start = False self._authority_release_safe = True self.state = PostgresState.STOPPED self.detail = ( 'PostgreSQL observer closed after the in-flight start was rejected; ' 'no controller-owned side effect was created or stopped' ) if result.detail: self.detail += f': {result.detail}' return True if operation == 'probe': self._handle_probe(result, now) elif operation == 'start': self._handle_start(result, now) elif operation == 'stop': if isinstance(result, StopResult) and result.completed and result.stopped: self._owned_start = False self._authority_release_safe = True if self._compensating_start_failure and not self._shutdown_requested: self._compensating_start_failure = False self._stop_submitted = False self._enter_backoff( now, result.detail or 'uncertain owned PostgreSQL start was safely compensated', ) else: self.state = PostgresState.STOPPED self.detail = result.detail or 'identity-verified PostgreSQL stop completed' else: self._compensating_start_failure = False self.state = PostgresState.STOP_FAILED self._authority_release_safe = False self.detail = result.detail if isinstance(result, StopResult) else 'invalid PostgreSQL stop result' return True def tick(self, now=None): now = time.monotonic() if now is None else float(now) self._consume_completion(now) if self.state in (PostgresState.DISABLED, PostgresState.STOPPED, PostgresState.STOP_FAILED): return self.state if self._shutdown_requested: self.state = PostgresState.STOPPING if self._future is None and not self._stop_submitted: self._submit('stop') return self.state if self._automatic_inhibited: return self.state if self._future is not None or now < self.next_action_at: return self.state if self.state == PostgresState.BACKOFF: self.state = PostgresState.VERIFYING if self.state in ( PostgresState.VERIFYING, PostgresState.STARTING, PostgresState.RECOVERING, PostgresState.STABILIZING, PostgresState.READY, PostgresState.FOREIGN_OR_CONFIG_ERROR, ): self._submit('probe') return self.state def request_stop(self): if self.state == PostgresState.DISABLED: return if self.state == PostgresState.STOPPED: return if self.state == PostgresState.STOP_FAILED: return if self._shutdown_requested: self.state = PostgresState.STOPPING return if self._operation == 'stop': self._shutdown_requested = True self._authority_release_safe = False self.state = PostgresState.STOPPING self.detail = 'coordinated shutdown is waiting for in-flight PostgreSQL compensation' return self._shutdown_requested = True self._authority_release_safe = False self._generation += 1 self.state = PostgresState.STOPPING self.detail = 'coordinated shutdown requested' def inhibit_lifecycle(self, detail): """Stop all automatic lifecycle submissions without taking process action.""" if self._shutdown_requested: return self._automatic_inhibited = True self._lifecycle_inert = True if self._operation != 'stop': self._generation += 1 self.state = PostgresState.FOREIGN_OR_CONFIG_ERROR self.detail = str(detail or 'automatic PostgreSQL lifecycle is inhibited') self.stable_since = None self._stable_epoch = '' self._ready_epoch = '' def retry_failed_stop(self): """Re-arm an identity-safe stop while endpoint authority is still held.""" if self.state != PostgresState.STOP_FAILED or self._future is not None: return False self._generation += 1 self._shutdown_requested = True self._stop_submitted = False self._authority_release_safe = False self.state = PostgresState.STOPPING self.detail = 'retrying identity-safe PostgreSQL compensation under retained authority' return True def snapshot(self): return { 'state': self.state.value, 'ready': self.ready, 'failures': self.failures, 'detail': self.detail, 'operation': self._operation, 'generation': self._generation, 'stop_succeeded': self.stop_succeeded, 'lifecycle_inert': self._lifecycle_inert, 'automatic_inhibited': self._automatic_inhibited, 'authority_release_safe': self.authority_release_safe, 'inflight_start': self.has_inflight_start, 'lifecycle_action_required': self.lifecycle_action_required, } def close(self, wait=False, timeout_sec=None): """Drain lifecycle work and return whether endpoint authority may be released.""" timeout = self.shutdown_timeout_sec if timeout_sec is None else max(0.0, float(timeout_sec)) deadline = time.monotonic() + timeout observer_only = ( self.state not in (PostgresState.DISABLED, PostgresState.STOPPED) and not self.lifecycle_action_required ) if observer_only: self._automatic_inhibited = True self._lifecycle_inert = True self._generation += 1 if self._future is not None: self._future.cancel() while self._future is not None and time.monotonic() <= deadline: self._consume_completion(time.monotonic()) if self._future is not None: time.sleep(0.01) if self._future is None: self.state = PostgresState.STOPPED self._authority_release_safe = True self.detail = ( 'observer-only PostgreSQL controller closed without stopping a cluster ' 'this controller did not start' ) else: if self.state == PostgresState.DISABLED: self._authority_release_safe = True elif self.state == PostgresState.STOP_FAILED: self.retry_failed_stop() elif self.state != PostgresState.STOPPED: self.request_stop() while not self.authority_release_safe and time.monotonic() <= deadline: self.tick() if self.state == PostgresState.STOP_FAILED and self._future is None: break time.sleep(0.01) safe = self.authority_release_safe if not safe: operation = self._operation or 'PostgreSQL compensation' self.state = PostgresState.STOP_FAILED self._authority_release_safe = False self.detail = ( f'bounded close left {operation} unresolved; endpoint authority release is unsafe' if self._future is not None else self.detail or 'PostgreSQL compensation did not confirm stopped state' ) return False if self._owns_executor: self._executor.shutdown(wait=True, cancel_futures=False) try: self.backend.close() except Exception: pass return True def controller_from_config(config, supervisor_config): return PostgresController( PostgresBackend( config, connect_timeout_sec=int(supervisor_config.get('postgres_connect_timeout_sec', DEFAULT_CONNECT_TIMEOUT_SEC) or DEFAULT_CONNECT_TIMEOUT_SEC), query_timeout_ms=int(supervisor_config.get('postgres_query_timeout_ms', DEFAULT_QUERY_TIMEOUT_MS) or DEFAULT_QUERY_TIMEOUT_MS), stop_timeout_sec=int(supervisor_config.get('postgres_stop_timeout_sec', 60) or 60), start_settle_timeout_sec=float(supervisor_config.get('postgres_start_settle_timeout_sec', 30) or 30), log_max_mb=int(supervisor_config.get('postgres_log_max_mb', 64) or 64), log_keep=int(supervisor_config.get('postgres_log_keep', 24) or 24), ), enabled=True, health_interval_sec=float(supervisor_config.get('postgres_health_interval_sec', 15) or 15), stable_ready_interval_sec=float(supervisor_config.get('postgres_stable_ready_sec', 60) or 60), ready_loss_grace_sec=float(supervisor_config.get('postgres_ready_loss_grace_sec', 45) or 45), backoff_base_sec=30, backoff_max_sec=600, shutdown_timeout_sec=float(supervisor_config.get('postgres_shutdown_timeout_sec', 120) or 120), ) def _load_config(path): try: import yaml except ImportError as exc: raise SystemExit('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 load_postgres_environment(config_path, config): global_config = (config or {}).get('global') or {} candidates = [] if global_config.get('root_dir'): candidates.append(os.path.join(global_config['root_dir'], '.env.postgres')) candidates.extend(( os.path.join(os.path.dirname(config_path), '..', '.env.postgres'), os.path.join(os.path.dirname(config_path), '.env.postgres'), )) loaded = None for candidate in candidates: candidate = os.path.abspath(candidate) if not os.path.isfile(candidate): continue with open(candidate, 'r', encoding='utf-8') as handle: for line in handle: text = line.strip() if not text or text.startswith('#') or '=' not in text: continue key, value = text.split('=', 1) key = key.strip() value = value.strip().strip('"').strip("'") if key and value and not os.getenv(key): os.environ[key] = value loaded = candidate break try: url = canonical_database_url() except DatabaseUrlError as exc: raise ClusterIdentityError(str(exc)) from exc if url: for key in list(os.environ): if key.upper().startswith('PG'): os.environ.pop(key, None) os.environ['SCANNER_DB_URL'] = url os.environ['DATABASE_URL'] = url return loaded _load_postgres_environment = load_postgres_environment @contextlib.contextmanager def _maintenance_shutdown_requests(): requested = False def request_shutdown(_signum, _frame): nonlocal requested requested = True previous = {} try: for signum in (signal.SIGINT, signal.SIGTERM): previous[signum] = signal.signal(signum, request_shutdown) yield lambda: requested finally: for signum, handler in previous.items(): signal.signal(signum, handler) def _owned_maintenance(config, *, start, shutdown_requested): """Caller holds cluster authority throughout this backend's ownership.""" backend = PostgresBackend(config) failure = None stopped = False result = None try: if start: result = maintenance_start(config, backend, shutdown_requested=shutdown_requested) else: result = maintenance_stop(config, backend) stopped = True if shutdown_requested(): raise ClusterIdentityError('maintenance PostgreSQL shutdown requested') label = 'authenticated-ready' if start else 'stopped' print(f'Maintenance PostgreSQL {label}: {result.detail}', flush=True) except BaseException as exc: failure = exc # Once compensation is required, mutable backend flags cannot replace a # completed StopResult (an interrupt may have lost that result). requires_stop = not start or result is not None or getattr(backend, 'owns_start', True) while True: try: if failure is None and shutdown_requested(): failure = ClusterIdentityError('maintenance PostgreSQL shutdown requested') if failure is not None and not stopped: # A rejected observational start owns no lifecycle. A stop, # or any possibly launched start, needs positive stop proof. if requires_stop: try: print('Maintenance PostgreSQL FAILED_HOLD: retaining backend and cluster authority until identity-verified stop.', flush=True) except BaseException: pass maintenance_stop(config, backend) stopped = True # READY is an intentional successful handoff to the next CLI. # Every failure path instead gets here only after safe compensation. backend.close() break except BaseException as exc: if failure is None: failure = exc try: time.sleep(1) except BaseException: pass if failure is not None: action = 'start' if start else 'stop' raise ClusterIdentityError( f'maintenance PostgreSQL {action} failed: {type(failure).__name__}: {failure}' ) from failure return result def maintenance_start(config, backend=None, *, shutdown_requested=None): """A supplied backend stays caller-owned, including failed/uncertain starts.""" if backend is None: with _maintenance_shutdown_requests() as requested: return _owned_maintenance(config, start=True, shutdown_requested=requested) if shutdown_requested is not None and shutdown_requested(): raise ClusterIdentityError('maintenance PostgreSQL shutdown requested') print('Maintenance PostgreSQL start request: validating offline cluster', flush=True) result = backend.start() if not result.accepted: raise ClusterIdentityError(result.detail or 'maintenance PostgreSQL start was refused') deadline = time.monotonic() + max(30.0, float(getattr(backend, 'start_settle_timeout_sec', 30.0)) + 10.0) while time.monotonic() < deadline: if shutdown_requested is not None and shutdown_requested(): raise ClusterIdentityError('maintenance PostgreSQL shutdown requested') probe = backend.probe() if probe.kind == ProbeKind.READY: return probe if probe.kind in (ProbeKind.STOPPED, ProbeKind.FOREIGN_OR_CONFIG_ERROR, ProbeKind.OWNED_START_UNCERTAIN): raise ClusterIdentityError(probe.detail or f'maintenance PostgreSQL entered {probe.kind.value}') time.sleep(0.25) raise ClusterIdentityError('maintenance PostgreSQL did not become authenticated-ready before timeout') def maintenance_stop(config, backend=None): """A supplied backend gets one bounded stop and is never implicitly closed.""" if backend is None: with _maintenance_shutdown_requests() as requested: return _owned_maintenance(config, start=False, shutdown_requested=requested) result = backend.stop() if not isinstance(result, StopResult) or result.completed is not True or result.stopped is not True: detail = result.detail if isinstance(result, StopResult) else '' raise ClusterIdentityError(detail or 'maintenance PostgreSQL stop did not complete') return result def main(): parser = argparse.ArgumentParser(description='Bootstrap or verify bundled PostgreSQL cluster authority while runtime sources are stopped.') parser.add_argument('action', choices=('initialize-empty', 'bootstrap', 'verify', 'maintenance-start', 'maintenance-stop')) parser.add_argument('--config', default=os.path.join(os.path.dirname(__file__), 'config.yaml')) args = parser.parse_args() config_path = os.path.abspath(args.config) config = _load_config(config_path) try: preflight_lifecycle_paths(config_path, config, authority_profile='server') load_postgres_environment(config_path, config) endpoint_dsn = canonical_database_url() if not endpoint_dsn: raise ClusterIdentityError('canonical managed PostgreSQL DSN is required') if args.action.startswith('maintenance-'): print('Maintenance PostgreSQL preflight: OK', flush=True) with _maintenance_shutdown_requests() as requested: with ClusterAuthorityLock(config, endpoint_dsn=endpoint_dsn): _owned_maintenance( config, start=args.action == 'maintenance-start', shutdown_requested=requested, ) return if args.action == 'initialize-empty': identity = initialize_empty(config) print(f'Initialized and bound stopped PostgreSQL cluster: {postgres_runtime_paths(config)["identity_path"]}') else: with ClusterAuthorityLock(config, endpoint_dsn=endpoint_dsn): if args.action == 'bootstrap': identity = bootstrap_cluster_identity(config) print(f'Wrote verified private cluster identity: {postgres_runtime_paths(config)["identity_path"]}') elif args.action == 'verify': identity = verify_cluster_identity(config) print(f'Verified private cluster identity: {postgres_runtime_paths(config)["identity_path"]}') except (OSError, ValueError, subprocess.SubprocessError) as exc: raise SystemExit(f'PostgreSQL cluster identity {args.action} failed closed: {exc}') from exc print(f'PostgreSQL {identity["pg_major"]} system_identifier={identity["system_identifier"]}') if __name__ == '__main__': main()