import ctypes import errno import hashlib import json import os import re import stat import threading from dataclasses import dataclass from enum import Enum MAX_PRIVATE_JSON_BYTES = 256 * 1024 MAX_EXTENDED_PRIVATE_JSON_BYTES = 64 * 1024 * 1024 _MACHINE_MUTEX_GUARD = threading.Lock() _MACHINE_MUTEX_NAMES = set() if os.name == 'nt': import msvcrt from ctypes import wintypes _ULONG_PTR = ctypes.c_ulonglong if ctypes.sizeof(ctypes.c_void_p) == 8 else wintypes.DWORD class _OVERLAPPED(ctypes.Structure): _fields_ = [ ('Internal', _ULONG_PTR), ('InternalHigh', _ULONG_PTR), ('Offset', wintypes.DWORD), ('OffsetHigh', wintypes.DWORD), ('hEvent', wintypes.HANDLE), ] class _SECURITY_ATTRIBUTES(ctypes.Structure): _fields_ = [ ('nLength', wintypes.DWORD), ('lpSecurityDescriptor', ctypes.c_void_p), ('bInheritHandle', wintypes.BOOL), ] class _SID_AND_ATTRIBUTES(ctypes.Structure): _fields_ = [('Sid', ctypes.c_void_p), ('Attributes', wintypes.DWORD)] class _TOKEN_USER(ctypes.Structure): _fields_ = [('User', _SID_AND_ATTRIBUTES)] class _FILETIME(ctypes.Structure): _fields_ = [('dwLowDateTime', wintypes.DWORD), ('dwHighDateTime', wintypes.DWORD)] class _BY_HANDLE_FILE_INFORMATION(ctypes.Structure): _fields_ = [ ('dwFileAttributes', wintypes.DWORD), ('ftCreationTime', _FILETIME), ('ftLastAccessTime', _FILETIME), ('ftLastWriteTime', _FILETIME), ('dwVolumeSerialNumber', wintypes.DWORD), ('nFileSizeHigh', wintypes.DWORD), ('nFileSizeLow', wintypes.DWORD), ('nNumberOfLinks', wintypes.DWORD), ('nFileIndexHigh', wintypes.DWORD), ('nFileIndexLow', wintypes.DWORD), ] _P_OVERLAPPED = ctypes.POINTER(_OVERLAPPED) _P_SECURITY_ATTRIBUTES = ctypes.POINTER(_SECURITY_ATTRIBUTES) _P_VOID_P = ctypes.POINTER(ctypes.c_void_p) _P_DWORD = ctypes.POINTER(wintypes.DWORD) _P_HANDLE = ctypes.POINTER(wintypes.HANDLE) _P_LPWSTR = ctypes.POINTER(wintypes.LPWSTR) _P_TOKEN_USER = ctypes.POINTER(_TOKEN_USER) _P_BY_HANDLE_FILE_INFORMATION = ctypes.POINTER(_BY_HANDLE_FILE_INFORMATION) _KERNEL32 = ctypes.WinDLL('kernel32', use_last_error=True) _ADVAPI32 = ctypes.WinDLL('advapi32', use_last_error=True) _SHELL32 = ctypes.WinDLL('shell32', use_last_error=True) _LOCK_FILE_EX = _KERNEL32.LockFileEx _LOCK_FILE_EX.argtypes = [ wintypes.HANDLE, wintypes.DWORD, wintypes.DWORD, wintypes.DWORD, wintypes.DWORD, _P_OVERLAPPED, ] _LOCK_FILE_EX.restype = wintypes.BOOL _UNLOCK_FILE_EX = _KERNEL32.UnlockFileEx _UNLOCK_FILE_EX.argtypes = [ wintypes.HANDLE, wintypes.DWORD, wintypes.DWORD, wintypes.DWORD, _P_OVERLAPPED, ] _UNLOCK_FILE_EX.restype = wintypes.BOOL _CREATE_MUTEX = _KERNEL32.CreateMutexW _CREATE_MUTEX.argtypes = [_P_SECURITY_ATTRIBUTES, wintypes.BOOL, wintypes.LPCWSTR] _CREATE_MUTEX.restype = wintypes.HANDLE _WAIT_FOR_SINGLE_OBJECT = _KERNEL32.WaitForSingleObject _WAIT_FOR_SINGLE_OBJECT.argtypes = [wintypes.HANDLE, wintypes.DWORD] _WAIT_FOR_SINGLE_OBJECT.restype = wintypes.DWORD _RELEASE_MUTEX = _KERNEL32.ReleaseMutex _RELEASE_MUTEX.argtypes = [wintypes.HANDLE] _RELEASE_MUTEX.restype = wintypes.BOOL _CLOSE_HANDLE = _KERNEL32.CloseHandle _CLOSE_HANDLE.argtypes = [wintypes.HANDLE] _CLOSE_HANDLE.restype = wintypes.BOOL _LOCAL_FREE = _KERNEL32.LocalFree _LOCAL_FREE.argtypes = [wintypes.HLOCAL] _LOCAL_FREE.restype = wintypes.HLOCAL _GET_CURRENT_PROCESS = _KERNEL32.GetCurrentProcess _GET_CURRENT_PROCESS.argtypes = [] _GET_CURRENT_PROCESS.restype = wintypes.HANDLE _CREATE_FILE = _KERNEL32.CreateFileW _CREATE_FILE.argtypes = [ wintypes.LPCWSTR, wintypes.DWORD, wintypes.DWORD, ctypes.c_void_p, wintypes.DWORD, wintypes.DWORD, wintypes.HANDLE, ] _CREATE_FILE.restype = wintypes.HANDLE _GET_FILE_INFORMATION = _KERNEL32.GetFileInformationByHandle _GET_FILE_INFORMATION.argtypes = [wintypes.HANDLE, _P_BY_HANDLE_FILE_INFORMATION] _GET_FILE_INFORMATION.restype = wintypes.BOOL _MOVE_FILE_EX = _KERNEL32.MoveFileExW _MOVE_FILE_EX.argtypes = [wintypes.LPCWSTR, wintypes.LPCWSTR, wintypes.DWORD] _MOVE_FILE_EX.restype = wintypes.BOOL _CONVERT_SDDL = _ADVAPI32.ConvertStringSecurityDescriptorToSecurityDescriptorW _CONVERT_SDDL.argtypes = [wintypes.LPCWSTR, wintypes.DWORD, _P_VOID_P, _P_DWORD] _CONVERT_SDDL.restype = wintypes.BOOL _SET_KERNEL_OBJECT_SECURITY = _ADVAPI32.SetKernelObjectSecurity _SET_KERNEL_OBJECT_SECURITY.argtypes = [wintypes.HANDLE, wintypes.DWORD, ctypes.c_void_p] _SET_KERNEL_OBJECT_SECURITY.restype = wintypes.BOOL _OPEN_PROCESS_TOKEN = _ADVAPI32.OpenProcessToken _OPEN_PROCESS_TOKEN.argtypes = [wintypes.HANDLE, wintypes.DWORD, _P_HANDLE] _OPEN_PROCESS_TOKEN.restype = wintypes.BOOL _GET_TOKEN_INFORMATION = _ADVAPI32.GetTokenInformation _GET_TOKEN_INFORMATION.argtypes = [ wintypes.HANDLE, wintypes.DWORD, ctypes.c_void_p, wintypes.DWORD, _P_DWORD, ] _GET_TOKEN_INFORMATION.restype = wintypes.BOOL _CONVERT_SID = _ADVAPI32.ConvertSidToStringSidW _CONVERT_SID.argtypes = [ctypes.c_void_p, _P_LPWSTR] _CONVERT_SID.restype = wintypes.BOOL _GET_FILE_SECURITY = _ADVAPI32.GetFileSecurityW _GET_FILE_SECURITY.argtypes = [ wintypes.LPCWSTR, wintypes.DWORD, ctypes.c_void_p, wintypes.DWORD, _P_DWORD, ] _GET_FILE_SECURITY.restype = wintypes.BOOL _SECURITY_DESCRIPTOR_TO_SDDL = _ADVAPI32.ConvertSecurityDescriptorToStringSecurityDescriptorW _SECURITY_DESCRIPTOR_TO_SDDL.argtypes = [ ctypes.c_void_p, wintypes.DWORD, wintypes.DWORD, _P_LPWSTR, _P_DWORD, ] _SECURITY_DESCRIPTOR_TO_SDDL.restype = wintypes.BOOL _SH_GET_FOLDER_PATH = _SHELL32.SHGetFolderPathW _SH_GET_FOLDER_PATH.argtypes = [ wintypes.HWND, ctypes.c_int, wintypes.HANDLE, wintypes.DWORD, wintypes.LPWSTR, ] _SH_GET_FOLDER_PATH.restype = ctypes.c_long else: _OVERLAPPED = _SECURITY_ATTRIBUTES = _SID_AND_ATTRIBUTES = None _TOKEN_USER = _FILETIME = _BY_HANDLE_FILE_INFORMATION = None _KERNEL32 = _ADVAPI32 = _SHELL32 = None class PrivateFileError(OSError): pass class PrivatePathState(str, Enum): PRESENT = 'present' ABSENT = 'absent' UNKNOWN = 'unknown' @dataclass(frozen=True) class PrivatePathInspection: state: PrivatePathState path: str detail: str = '' stat_result: object = None def inspect_private_relative_path(root, relative_path): """Inspect an indexed path without treating inaccessible storage as absence.""" try: root = require_private_directory(os.path.abspath(root), create=False) relative = str(relative_path or '').replace('/', os.sep) if not relative or os.path.isabs(relative): raise ValueError('private relative path is empty or absolute') path = os.path.abspath(os.path.join(root, relative)) if path == root or os.path.commonpath((root, path)) != root: raise ValueError('private relative path escapes its root') parent = os.path.dirname(path) require_private_directory(parent, create=False) # Opening the shard proves traversal/access independently of target lstat. with os.scandir(parent): pass try: details = os.lstat(path) except FileNotFoundError: return PrivatePathInspection(PrivatePathState.ABSENT, path) except OSError as exc: return PrivatePathInspection( PrivatePathState.UNKNOWN, path, f'{type(exc).__name__}: {exc}', ) return PrivatePathInspection(PrivatePathState.PRESENT, path, stat_result=details) except (OSError, ValueError) as exc: candidate = os.path.abspath(os.path.join(os.path.abspath(root), str(relative_path or ''))) return PrivatePathInspection( PrivatePathState.UNKNOWN, candidate, f'{type(exc).__name__}: {exc}', ) class PrivateFileLock: """Cross-process lifetime lock backed by one byte of a private file.""" def __init__(self, path): self.path = os.path.normcase(os.path.abspath(os.fspath(path))) self._file = None self._overlapped = None self._acquired = False @property def acquired(self): return self._acquired def acquire(self): if self._acquired: return self parent = os.path.dirname(self.path) require_private_directory(parent, create=False) reject_reparse_components(self.path) flags = os.O_RDWR | os.O_CREAT if hasattr(os, 'O_BINARY'): flags |= os.O_BINARY if hasattr(os, 'O_NOFOLLOW'): flags |= os.O_NOFOLLOW existed = os.path.lexists(self.path) if existed: descriptor = os.open(self.path, flags, 0o600) else: try: descriptor = os.open(self.path, flags | os.O_EXCL, 0o600) except FileExistsError: existed = True reject_reparse_components(self.path) descriptor = os.open(self.path, flags, 0o600) try: if existed: if not private_file_ready(self.path): raise PrivateFileError(f'private lock file ACL is not ready: {self.path}') else: harden_private_file(self.path) file_handle = os.fdopen(descriptor, 'r+b', buffering=0) descriptor = None if os.name == 'nt': overlapped = _OVERLAPPED() handle = wintypes.HANDLE(msvcrt.get_osfhandle(file_handle.fileno())) flags_value = 0x00000002 | 0x00000001 # EXCLUSIVE_LOCK | FAIL_IMMEDIATELY if not _LOCK_FILE_EX(handle, flags_value, 0, 1, 0, ctypes.byref(overlapped)): error = ctypes.get_last_error() file_handle.close() raise BlockingIOError(error, 'private lifecycle lock is already held', self.path) self._overlapped = overlapped else: import fcntl try: fcntl.flock(file_handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) except OSError as exc: file_handle.close() raise BlockingIOError(exc.errno, 'private lifecycle lock is already held', self.path) from exc self._file = file_handle self._acquired = True return self finally: if descriptor is not None: os.close(descriptor) def release(self): if not self._acquired: return self._acquired = False file_handle = self._file self._file = None try: if os.name == 'nt': handle = wintypes.HANDLE(msvcrt.get_osfhandle(file_handle.fileno())) _UNLOCK_FILE_EX(handle, 0, 1, 0, ctypes.byref(self._overlapped)) self._overlapped = None else: import fcntl fcntl.flock(file_handle.fileno(), fcntl.LOCK_UN) finally: file_handle.close() def __enter__(self): return self.acquire() def __exit__(self, exc_type, value, traceback): self.release() def canonical_cluster_data_directory(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') from paths import resolve_postgres_data_dir data_dir = reject_reparse_components( resolve_postgres_data_dir(global_config, runtime_dir), ) drive, _ = os.path.splitdrive(data_dir) root = drive + os.sep if drive else os.sep if os.path.normcase(os.path.normpath(data_dir)) == os.path.normcase(os.path.normpath(root)): raise PrivateFileError('PostgreSQL data directory cannot be a filesystem or volume root') return canonical_path(data_dir) def cluster_authority_lock_path(config): data_dir = canonical_cluster_data_directory(config) digest = hashlib.sha256(data_dir.encode('utf-8', errors='strict')).hexdigest() 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') return os.path.join(canonical_path(os.path.join(runtime_dir, 'postgres')), f'.cluster-authority-{digest}.lock') def _postgres_environment_values(config): keys = { 'TRUF_MANAGED_POSTGRES_DSN', 'SCANNER_DB_URL', 'DATABASE_URL', 'TRUF_POSTGRES_DB', 'TRUF_POSTGRES_USER', 'TRUF_POSTGRES_PORT', } values = {key: os.getenv(key) for key in keys if os.getenv(key)} global_config = (config or {}).get('global') or {} candidates = [] for parent in (global_config.get('root_dir'), global_config.get('project_dir')): if parent: candidates.append(os.path.join(parent, '.env.postgres')) seen = set() for candidate in candidates: candidate = os.path.abspath(candidate) normalized = os.path.normcase(candidate) if normalized in seen or not os.path.isfile(candidate): continue seen.add(normalized) try: 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 in keys and value and key not in values: values[key] = value except OSError as exc: raise PrivateFileError(f'unable to determine PostgreSQL endpoint authority from {candidate}') from exc if any(values.get(key) for key in ('TRUF_MANAGED_POSTGRES_DSN', 'SCANNER_DB_URL', 'DATABASE_URL')): break return values def canonical_cluster_endpoint_identity(config=None, database_url=None): """Return the credential-free identity of the PostgreSQL mutation endpoint.""" if database_url is None: raise PrivateFileError('canonical managed PostgreSQL DSN must be passed explicitly for endpoint authority') from db_backend import POSTGRES_APPLICATION_SCHEMA, parse_postgres_url parsed = parse_postgres_url(database_url) host = parsed['host'] port = parsed['port'] database = parsed['database'] host = str(host or '').strip().lower().rstrip('.') if host in ('localhost', '::1') or re.fullmatch(r'127(?:\.\d{1,3}){3}', host): host = '127.0.0.1' return json.dumps( {'database': database, 'host': host, 'port': int(port), 'schema': POSTGRES_APPLICATION_SCHEMA}, ensure_ascii=True, sort_keys=True, separators=(',', ':'), ) def _windows_common_appdata(): buffer = ctypes.create_unicode_buffer(32768) result = _SH_GET_FOLDER_PATH(None, 0x0023, None, 0, buffer) # CSIDL_COMMON_APPDATA if result != 0 or not buffer.value: raise PrivateFileError('unable to resolve the machine authority directory') return os.path.abspath(buffer.value) def _cluster_endpoint_lock_root(): if os.name == 'nt': return os.path.join(_windows_common_appdata(), 'Truf', 'authority') return '/run/truf/authority' def cluster_endpoint_authority_lock_path(config=None, endpoint_identity=None, endpoint_dsn=None): endpoint = endpoint_identity or canonical_cluster_endpoint_identity(config, database_url=endpoint_dsn) digest = hashlib.sha256(endpoint.encode('utf-8', errors='strict')).hexdigest() return os.path.join(_cluster_endpoint_lock_root(), f'endpoint-{digest}.lock') def cluster_endpoint_mutex_name(config=None, endpoint_identity=None, endpoint_dsn=None): endpoint = endpoint_identity or canonical_cluster_endpoint_identity(config, database_url=endpoint_dsn) digest = hashlib.sha256(endpoint.encode('utf-8', errors='strict')).hexdigest() return f'Global\\Truf.ClusterEndpoint.{digest}' class _WindowsGlobalMutex: """Non-recursive process wrapper around a private cross-session mutex.""" def __init__(self, name): self.name = str(name) self._handle = None self._acquired = False @property def acquired(self): return self._acquired def acquire(self): if self._acquired: return self with _MACHINE_MUTEX_GUARD: if self.name in _MACHINE_MUTEX_NAMES: raise BlockingIOError(0, 'private machine endpoint authority is already held', self.name) _MACHINE_MUTEX_NAMES.add(self.name) descriptor = ctypes.c_void_p() sid = _windows_current_user_sid() sddl = f'D:P(A;;GA;;;{sid})(A;;GA;;;SY)(A;;GA;;;BA)' try: if not _CONVERT_SDDL(sddl, 1, ctypes.byref(descriptor), None): raise ctypes.WinError(ctypes.get_last_error()) attributes = _SECURITY_ATTRIBUTES(ctypes.sizeof(_SECURITY_ATTRIBUTES), descriptor, False) handle = _CREATE_MUTEX(ctypes.byref(attributes), False, self.name) if not handle: raise ctypes.WinError(ctypes.get_last_error()) result = _WAIT_FOR_SINGLE_OBJECT(handle, 0) if result not in (0, 0x00000080): # WAIT_OBJECT_0 / WAIT_ABANDONED _CLOSE_HANDLE(handle) if result == 0x00000102: # WAIT_TIMEOUT raise BlockingIOError(0, 'private machine endpoint authority is already held', self.name) raise ctypes.WinError(ctypes.get_last_error()) self._handle = handle self._acquired = True return self except BaseException: with _MACHINE_MUTEX_GUARD: _MACHINE_MUTEX_NAMES.discard(self.name) raise finally: if descriptor: _LOCAL_FREE(descriptor) def release(self): if not self._acquired: return handle = self._handle self._handle = None self._acquired = False try: if handle and not _RELEASE_MUTEX(handle): raise ctypes.WinError(ctypes.get_last_error()) finally: if handle: _CLOSE_HANDLE(handle) with _MACHINE_MUTEX_GUARD: _MACHINE_MUTEX_NAMES.discard(self.name) class ClusterAuthorityLock: """Cross-process authority over both cluster storage and its SQL endpoint.""" def __init__(self, config, create_parent=False, endpoint_dsn=None): if endpoint_dsn is None: raise PrivateFileError('cluster endpoint authority requires the caller-selected canonical DSN') self.data_directory = canonical_cluster_data_directory(config) self.path = cluster_authority_lock_path(config) self.endpoint_identity = canonical_cluster_endpoint_identity(config, database_url=endpoint_dsn) self.endpoint_path = cluster_endpoint_authority_lock_path(config, self.endpoint_identity) self.endpoint_mutex_name = cluster_endpoint_mutex_name(endpoint_identity=self.endpoint_identity) self.paths = (self.endpoint_path, self.path) self.create_parent = bool(create_parent) self._locks = [] @property def acquired(self): return len(self._locks) == len(self.paths) and all(lock.acquired for lock in self._locks) def acquire(self): data_parent = os.path.dirname(self.path) if self.create_parent: reject_reparse_components(os.path.dirname(data_parent)) os.makedirs(data_parent, mode=0o700, exist_ok=True) reject_reparse_components(data_parent) harden_private_directory(data_parent) try: if os.name == 'nt': self._locks.append(_WindowsGlobalMutex(self.endpoint_mutex_name).acquire()) else: endpoint_parent = os.path.dirname(self.endpoint_path) if not os.path.exists(endpoint_parent): os.makedirs(endpoint_parent, mode=0o700, exist_ok=False) harden_private_directory(endpoint_parent) require_private_directory(endpoint_parent, create=False) self._locks.append(PrivateFileLock(self.endpoint_path).acquire()) self._locks.append(PrivateFileLock(self.path).acquire()) except BlockingIOError as exc: self.release() raise BlockingIOError( getattr(exc, 'errno', 0), 'another runtime or maintenance process owns this PostgreSQL data directory or endpoint authority', getattr(exc, 'filename', None) or self.path, ) from exc except BaseException: self.release() raise return self def release(self): while self._locks: self._locks.pop().release() def __enter__(self): return self.acquire() def __exit__(self, exc_type, value, traceback): self.release() def canonical_path(path): return os.path.normcase(os.path.realpath(os.path.abspath(os.fspath(path)))) def sha256_file(path): digest = hashlib.sha256() with open(path, 'rb') as handle: while True: block = handle.read(1024 * 1024) if not block: break digest.update(block) return digest.hexdigest() def is_reparse_point(path): try: info = os.lstat(path) except OSError: return False if stat.S_ISLNK(info.st_mode): return True return bool(getattr(info, 'st_file_attributes', 0) & 0x00000400) def reject_reparse_components(path): absolute = os.path.abspath(os.fspath(path)) drive, tail = os.path.splitdrive(absolute) current = drive + os.sep if drive else os.sep for part in tail.strip(os.sep).split(os.sep): if not part: continue current = os.path.join(current, part) if os.path.lexists(current) and is_reparse_point(current): raise PrivateFileError(f'private path contains a link or reparse point: {current}') return absolute def _windows_current_user_sid(): token = wintypes.HANDLE() if not _OPEN_PROCESS_TOKEN(_GET_CURRENT_PROCESS(), 0x0008, ctypes.byref(token)): raise ctypes.WinError(ctypes.get_last_error()) try: needed = wintypes.DWORD() _GET_TOKEN_INFORMATION(token, 1, None, 0, ctypes.byref(needed)) if not needed.value: raise ctypes.WinError(ctypes.get_last_error()) buffer = ctypes.create_string_buffer(needed.value) if not _GET_TOKEN_INFORMATION(token, 1, buffer, needed.value, ctypes.byref(needed)): raise ctypes.WinError(ctypes.get_last_error()) user = ctypes.cast(buffer, _P_TOKEN_USER).contents output = wintypes.LPWSTR() if not _CONVERT_SID(user.User.Sid, ctypes.byref(output)): raise ctypes.WinError(ctypes.get_last_error()) try: return output.value finally: _LOCAL_FREE(output) finally: _CLOSE_HANDLE(token) def _windows_private_sddl(path): security_information = 0x00000001 | 0x00000004 needed = wintypes.DWORD() _GET_FILE_SECURITY(path, security_information, None, 0, ctypes.byref(needed)) error = ctypes.get_last_error() if not needed.value or error not in (0, 122): raise ctypes.WinError(error) buffer = ctypes.create_string_buffer(needed.value) if not _GET_FILE_SECURITY(path, security_information, buffer, needed.value, ctypes.byref(needed)): raise ctypes.WinError(ctypes.get_last_error()) output = wintypes.LPWSTR() if not _SECURITY_DESCRIPTOR_TO_SDDL(buffer, 1, security_information, ctypes.byref(output), None): raise ctypes.WinError(ctypes.get_last_error()) try: return output.value or '' finally: _LOCAL_FREE(output) def _harden_windows_path(path, directory=False): # Protected DACL: the exact current owner SID, LocalSystem, and local Administrators only. # Directories propagate the same protected allowlist to newly-created children. descriptor = ctypes.c_void_p() ace_flags = 'OICI' if directory else '' owner_sid = _windows_current_user_sid() current_sddl = _windows_private_sddl(path).upper() owner_match = re.search(r'O:([^:()]+?)(?=[GDS]:|$)', current_sddl) if not owner_match or owner_match.group(1) != owner_sid.upper(): raise PrivateFileError(f'refusing to harden a path not owned by the current user: {path}') trustees = [] for trustee in (owner_sid, 'SY', 'BA'): if trustee not in trustees: trustees.append(trustee) sddl = 'D:P' + ''.join(f'(A;{ace_flags};FA;;;{trustee})' for trustee in trustees) if not _CONVERT_SDDL(sddl, 1, ctypes.byref(descriptor), None): raise ctypes.WinError(ctypes.get_last_error()) access = 0x00020000 | 0x00040000 # READ_CONTROL | WRITE_DAC sharing = 0x00000001 | 0x00000002 | 0x00000004 flags = 0x00200000 | 0x02000000 # OPEN_REPARSE_POINT | BACKUP_SEMANTICS handle = _CREATE_FILE(path, access, sharing, None, 3, flags, None) if handle == wintypes.HANDLE(-1).value: _LOCAL_FREE(descriptor) raise ctypes.WinError(ctypes.get_last_error()) try: information = _BY_HANDLE_FILE_INFORMATION() if not _GET_FILE_INFORMATION(handle, ctypes.byref(information)): raise ctypes.WinError(ctypes.get_last_error()) if information.dwFileAttributes & 0x00000400: raise PrivateFileError(f'refusing to harden a reparse point: {path}') security_information = 0x00000004 | 0x80000000 # DACL | PROTECTED_DACL if not _SET_KERNEL_OBJECT_SECURITY(handle, security_information, descriptor): raise ctypes.WinError(ctypes.get_last_error()) finally: _CLOSE_HANDLE(handle) _LOCAL_FREE(descriptor) def harden_private_file(path): reject_reparse_components(path) if os.name == 'nt': _harden_windows_path(path, directory=False) else: os.chmod(path, stat.S_IRUSR | stat.S_IWUSR) if not private_file_ready(path): raise PrivateFileError(f'private-file ACL verification failed: {path}') def harden_private_directory(path): reject_reparse_components(path) if os.name == 'nt': _harden_windows_path(path, directory=True) else: os.chmod(path, stat.S_IRWXU) if not private_directory_ready(path): raise PrivateFileError(f'private-directory ACL verification failed: {path}') def ensure_private_directory(path, reject_reparse=False): if reject_reparse: reject_reparse_components(os.path.dirname(os.path.abspath(path)) or path) os.makedirs(path, mode=0o700, exist_ok=True) if reject_reparse: reject_reparse_components(path) harden_private_directory(path) return path def _private_acl_ready(path, directory=False): if os.name != 'nt': info = os.stat(path, follow_symlinks=False) owner_ok = not hasattr(os, 'geteuid') or info.st_uid == os.geteuid() return owner_ok and stat.S_IMODE(info.st_mode) & 0o077 == 0 reject_reparse_components(path) sddl = _windows_private_sddl(path).upper() owner_match = re.search(r'O:([^:()]+?)(?=[GDS]:|$)', sddl) if not owner_match or 'D:P' not in sddl: return False aliases = {'S-1-3-4': 'OW', 'S-1-5-18': 'SY', 'S-1-5-32-544': 'BA'} owner = aliases.get(owner_match.group(1), owner_match.group(1)) current_owner = _windows_current_user_sid().upper() current_owner = aliases.get(current_owner, current_owner) if owner != current_owner: return False aces = re.findall(r'\(([^()]*)\)', sddl) trustees = set() expected_trustees = {current_owner, 'SY', 'BA'} if len(aces) != len(expected_trustees): return False expected_flags = {'OI', 'CI'} if directory else set() for ace in aces: fields = ace.split(';') flags = set(re.findall(r'OI|CI|IO|NP|ID', fields[1])) if len(fields) == 6 else set() if ( len(fields) != 6 or fields[0] != 'A' or flags != expected_flags or fields[2] != 'FA' or fields[3] or fields[4] ): return False trustees.add(aliases.get(fields[5], fields[5])) return trustees == expected_trustees def private_file_ready(path): try: reject_reparse_components(path) if not stat.S_ISREG(os.stat(path, follow_symlinks=False).st_mode): return False return _private_acl_ready(path, directory=False) except (OSError, ValueError): return False def private_directory_ready(path): try: reject_reparse_components(path) return stat.S_ISDIR(os.stat(path, follow_symlinks=False).st_mode) and _private_acl_ready(path, directory=True) except (OSError, ValueError): return False def read_stable_root_file(path, max_bytes, trusted_root): """Read one immutable, public root-owned file without following links.""" if os.name == 'nt' or type(max_bytes) is not int or max_bytes < 1: raise PrivateFileError('stable root file is unavailable') descriptor = None directories = [] try: absolute = os.path.abspath(os.fsdecode(path)) root = os.path.abspath(os.fsdecode(trusted_root)) if os.path.commonpath((absolute, root)) != root or absolute == root: raise PrivateFileError('stable root file is outside its authority') relative = os.path.relpath(absolute, root) parts = relative.split(os.sep) if not parts or any(part in ('', '.', '..') for part in parts): raise PrivateFileError('stable root file path is invalid') directory_flags = ( os.O_RDONLY | getattr(os, 'O_CLOEXEC', 0) | getattr(os, 'O_DIRECTORY', 0) | getattr(os, 'O_NOFOLLOW', 0) ) directory = os.open(root, directory_flags) directories.append(directory) for part in parts[:-1]: details = os.fstat(directory) if ( not stat.S_ISDIR(details.st_mode) or details.st_uid != 0 or details.st_gid != 0 or stat.S_IMODE(details.st_mode) & 0o022 ): raise PrivateFileError('stable root directory metadata is invalid') directory = os.open(part, directory_flags, dir_fd=directory) directories.append(directory) details = os.fstat(directory) if ( not stat.S_ISDIR(details.st_mode) or details.st_uid != 0 or details.st_gid != 0 or stat.S_IMODE(details.st_mode) & 0o022 ): raise PrivateFileError('stable root directory metadata is invalid') flags = os.O_RDONLY | getattr(os, 'O_CLOEXEC', 0) | getattr(os, 'O_NOFOLLOW', 0) descriptor = os.open(parts[-1], flags, dir_fd=directory) before = os.fstat(descriptor) def identity(details): return ( details.st_dev, details.st_ino, details.st_size, details.st_uid, details.st_gid, stat.S_IMODE(details.st_mode), details.st_nlink, getattr(details, 'st_mtime_ns', None), getattr(details, 'st_ctime_ns', None), ) if ( not stat.S_ISREG(before.st_mode) or before.st_nlink != 1 or before.st_uid != 0 or before.st_gid != 0 or stat.S_IMODE(before.st_mode) != 0o644 ): raise PrivateFileError('stable root file metadata is invalid') with os.fdopen(descriptor, 'rb') as handle: descriptor = None payload = handle.read(max_bytes + 1) after = os.fstat(handle.fileno()) current = os.stat(parts[-1], dir_fd=directory, follow_symlinks=False) if ( len(payload) > max_bytes or identity(before) != identity(after) or identity(after) != identity(current) ): raise PrivateFileError('stable root file changed while being read') return payload except PrivateFileError: raise except (OSError, TypeError, ValueError) as exc: raise PrivateFileError('stable root file is not ready') from exc finally: if descriptor is not None: os.close(descriptor) for directory in reversed(directories): os.close(directory) def fsync_directory(path): if os.name == 'nt' or not path: return flags = os.O_RDONLY if hasattr(os, 'O_DIRECTORY'): flags |= os.O_DIRECTORY descriptor = os.open(path, flags) try: os.fsync(descriptor) finally: os.close(descriptor) def durable_replace(source, destination): source = os.path.abspath(source) destination = os.path.abspath(destination) if os.path.dirname(source) != os.path.dirname(destination): raise ValueError('durable replacement must remain in the same directory') durable_move(source, destination) def durable_move(source, destination): source = os.path.abspath(source) destination = os.path.abspath(destination) reject_reparse_components(source) reject_reparse_components(os.path.dirname(destination)) if os.path.lexists(destination): reject_reparse_components(destination) source_parent = os.path.dirname(source) destination_parent = os.path.dirname(destination) if os.name == 'nt': movefile_replace_existing = 0x00000001 movefile_write_through = 0x00000008 if not _MOVE_FILE_EX(source, destination, movefile_replace_existing | movefile_write_through): raise ctypes.WinError(ctypes.get_last_error()) else: os.replace(source, destination) fsync_directory(destination_parent) if source_parent != destination_parent: fsync_directory(source_parent) def durable_publish(source, destination): """Atomically publish a new name without replacing an existing object.""" source = os.path.abspath(source) destination = os.path.abspath(destination) reject_reparse_components(source) reject_reparse_components(os.path.dirname(destination)) if os.path.lexists(destination): raise FileExistsError(destination) source_parent = os.path.dirname(source) destination_parent = os.path.dirname(destination) if os.path.splitdrive(source)[0].lower() != os.path.splitdrive(destination)[0].lower(): raise ValueError('durable publication must remain on one volume') if os.name == 'nt': movefile_write_through = 0x00000008 if not _MOVE_FILE_EX(source, destination, movefile_write_through): error = ctypes.get_last_error() if error in (80, 183): raise FileExistsError(destination) raise ctypes.WinError(error) else: os.link(source, destination, follow_symlinks=False) try: os.unlink(source) except BaseException: try: os.unlink(destination) except OSError: pass raise fsync_directory(destination_parent) if source_parent != destination_parent: fsync_directory(source_parent) def durable_publish_directory(source, destination): """Atomically move one exact directory to a new, non-existing name.""" source = os.path.abspath(source) destination = os.path.abspath(destination) source_parent = os.path.dirname(source) destination_parent = os.path.dirname(destination) if source == destination or not source_parent or not destination_parent: raise ValueError('durable directory publication paths are invalid') reject_reparse_components(source) reject_reparse_components(destination_parent) source_details = os.stat(source, follow_symlinks=False) if not stat.S_ISDIR(source_details.st_mode) or is_reparse_point(source): raise ValueError('durable directory publication source is not an exact directory') source_parent_details = os.stat(source_parent, follow_symlinks=False) destination_parent_details = os.stat(destination_parent, follow_symlinks=False) if source_parent_details.st_dev != destination_parent_details.st_dev: raise OSError(errno.EXDEV, 'durable directory publication cannot cross devices') if os.path.lexists(destination): raise FileExistsError(destination) if os.name == 'nt': movefile_write_through = 0x00000008 if not _MOVE_FILE_EX(source, destination, movefile_write_through): error = ctypes.get_last_error() if error in (80, 183): raise FileExistsError(destination) raise ctypes.WinError(error) return library = ctypes.CDLL(None, use_errno=True) renameat2 = getattr(library, 'renameat2', None) if renameat2 is None: raise OSError( errno.ENOSYS, 'atomic non-replacing directory publication requires renameat2', ) renameat2.argtypes = [ ctypes.c_int, ctypes.c_char_p, ctypes.c_int, ctypes.c_char_p, ctypes.c_uint, ] renameat2.restype = ctypes.c_int if renameat2( -100, os.fsencode(source), -100, os.fsencode(destination), 1, ) != 0: error = ctypes.get_errno() if error == errno.EEXIST: raise FileExistsError(destination) raise OSError(error, os.strerror(error), destination) fsync_directory(destination_parent) if source_parent != destination_parent: fsync_directory(source_parent) def durable_unlink(path): path = os.path.abspath(path) reject_reparse_components(path) parent = os.path.dirname(path) os.remove(path) fsync_directory(parent) def require_private_directory(path, create=False): """Require an exact private directory without silently repairing an existing ACL.""" absolute = reject_reparse_components(path) if not os.path.exists(absolute): if not create: raise PrivateFileError(f'private directory is absent: {absolute}') os.makedirs(absolute, mode=0o700, exist_ok=False) harden_private_directory(absolute) if not private_directory_ready(absolute): raise PrivateFileError(f'private directory ACL is not ready: {absolute}') return absolute def require_private_file(path): """Require an existing exact private regular file without changing it.""" absolute = reject_reparse_components(path) if not private_file_ready(absolute): raise PrivateFileError(f'private file ACL or owner is not ready: {absolute}') return absolute def require_protected_sensitive_file_parent(path): """Accept private runtime directories or protected root-owned read-only authorities.""" absolute = reject_reparse_components(path) if private_directory_ready(absolute): return absolute if os.name != 'nt': try: details = os.stat(absolute, follow_symlinks=False) mode = stat.S_IMODE(details.st_mode) effective_uid = os.geteuid() effective_gid = os.getegid() groups = set(os.getgroups()) if effective_uid == details.st_uid: search_bit = stat.S_IXUSR elif details.st_gid == effective_gid or details.st_gid in groups: search_bit = stat.S_IXGRP else: search_bit = stat.S_IXOTH if ( stat.S_ISDIR(details.st_mode) and details.st_uid == 0 and details.st_gid == 0 and mode & 0o022 == 0 and mode & search_bit ): return absolute except (OSError, ValueError): pass raise PrivateFileError(f'sensitive file parent ACL is not ready: {absolute}') def require_trusted_native_executable(path): """Return canonical native code trusted by the effective runtime user, without repairs. POSIX requires immutable root-owned code and parents; Windows retains the private-file policy. Verification failures raise PrivateFileError. """ try: if os.name == 'nt': return canonical_path(require_private_file(path)) raw_path = os.fsdecode(path) if not os.path.isabs(raw_path) or '\x00' in raw_path: raise PrivateFileError(f'trusted native executable must be an absolute path: {raw_path}') # Inspect before normalization so a link followed by /.. cannot disappear. current = raw_path while True: details = os.lstat(current) if stat.S_ISLNK(details.st_mode): raise PrivateFileError(f'trusted native path contains a link: {current}') if current == raw_path: if not stat.S_ISREG(details.st_mode): raise PrivateFileError(f'trusted native executable is not a regular file: {current}') elif not stat.S_ISDIR(details.st_mode): raise PrivateFileError(f'trusted native parent is not a directory: {current}') if details.st_uid != 0: raise PrivateFileError(f'trusted native path is not root-owned: {current}') if details.st_mode & (stat.S_ISUID | stat.S_ISGID | stat.S_IWGRP | stat.S_IWOTH): raise PrivateFileError(f'trusted native path has unsafe permissions: {current}') if os.access(current, os.W_OK, effective_ids=True): raise PrivateFileError(f'trusted native path is writable by the runtime user: {current}') if not os.access(current, os.X_OK, effective_ids=True): raise PrivateFileError(f'trusted native path is not executable/searchable by the runtime user: {current}') parent = os.path.dirname(current) if parent == current: break current = parent return canonical_path(raw_path) except PrivateFileError: raise except (OSError, TypeError, ValueError, NotImplementedError) as exc: raise PrivateFileError(f'trusted native executable verification failed: {path}: {exc}') from exc def preflight_lifecycle_paths(config_path, config, *, authority_profile='full'): """Read-only verification required before lifecycle secrets or logs are opened.""" if authority_profile not in ('full', 'discovery-producer', 'server'): raise PrivateFileError('unsupported lifecycle authority profile') discovery_producer = authority_profile in ('discovery-producer', 'server') global_config = (config or {}).get('global') or {} supervisor_config = (config or {}).get('supervisor') or {} runtime_dir = global_config.get('runtime_dir') directory_values = [ global_config.get('root_dir'), global_config.get('project_dir'), runtime_dir, global_config.get('control_dir'), supervisor_config.get('control_dir'), global_config.get('log_dir'), supervisor_config.get('log_dir'), global_config.get('work_dir'), global_config.get('results_dir'), global_config.get('queue_dir'), global_config.get('state_dir'), global_config.get('keycheck_dir'), global_config.get('postman_cache_dir'), global_config.get('gharchive_cache_dir'), global_config.get('result_spool_dir'), global_config.get('result_bundle_dir'), os.path.join(global_config.get('result_bundle_dir'), 'tmp') if global_config.get('result_bundle_dir') else None, os.path.join(global_config.get('result_bundle_dir'), 'ready') if global_config.get('result_bundle_dir') else None, os.path.join(global_config.get('result_bundle_dir'), 'quarantine') if global_config.get('result_bundle_dir') else None, os.path.join(runtime_dir, 'postgres') if runtime_dir else None, canonical_cluster_data_directory(config) if global_config.get('postgres_data_dir') else None, _cluster_endpoint_lock_root() if os.name != 'nt' else None, ] seen = set() try: require_private_file(config_path) for value in directory_values: if not value: continue absolute = os.path.abspath(value) normalized = os.path.normcase(absolute) if normalized in seen: continue seen.add(normalized) require_private_directory(absolute, create=False) sensitive_files = [ global_config.get('secrets_file'), global_config.get('proxy_file'), global_config.get('api_proxy_file'), global_config.get('download_proxy_file'), supervisor_config.get('instance_file'), supervisor_config.get('lock_file'), supervisor_config.get('supervisor_log'), supervisor_config.get('status_file'), supervisor_config.get('dashboard_log'), ] if not discovery_producer: sensitive_files.append(global_config.get('trufflehog_config')) policy_paths = [] if discovery_producer else [global_config.get('trufflehog_config')] from paths import resolve_optional_path if not discovery_producer: for source in ((config or {}).get('sources') or {}).values(): if isinstance(source, dict) and source.get('trufflehog_config'): policy_paths.append(resolve_optional_path(source['trufflehog_config'], global_config)) sensitive_files.extend(policy_paths) root_dir = global_config.get('root_dir') project_dir = global_config.get('project_dir') config_dir = os.path.dirname(os.path.abspath(config_path)) for parent in (root_dir, project_dir, config_dir, os.path.dirname(config_dir)): if parent: sensitive_files.append(os.path.join(parent, '.env.postgres')) for value in sensitive_files: if value and os.path.lexists(value): require_private_file(value) # Import lazily: lifecycle authority itself depends on this module. from lifecycle_authority import ( GIT_MANIFEST_NAME, TRUFFLEHOG_MANIFEST_NAME, manifest_authority_paths, resolve_manifest_executable, ) for path in manifest_authority_paths( global_config.get('project_dir'), global_config.get('trufflehog_path'), policy_paths=policy_paths, existing_only=True, include_executables=False, ): require_private_file(path) if authority_profile != 'discovery-producer': executable_values = [(GIT_MANIFEST_NAME, None)] if authority_profile == 'full': executable_values.insert( 0, (TRUFFLEHOG_MANIFEST_NAME, global_config.get('trufflehog_path')), ) for name, value in executable_values: require_trusted_native_executable(resolve_manifest_executable( value, name=name, app_dir=project_dir, )) except (OSError, ValueError) as exc: raise PrivateFileError( f'lifecycle sensitive-path preflight failed: {exc}. ' 'Run the offline hardening command before startup; runtime startup will not repair ACLs.' ) from exc return True def harden_private_tree(path): """Offline recursive hardening that never traverses a reparse point.""" absolute = reject_reparse_components(path) if not os.path.lexists(absolute): return 0 if is_reparse_point(absolute): raise PrivateFileError(f'refusing to harden a link or reparse point: {absolute}') if os.path.isfile(absolute): harden_private_file(absolute) return 1 if not os.path.isdir(absolute): raise PrivateFileError(f'unsupported private-tree entry: {absolute}') count = 0 stack = [absolute] directories = [] while stack: current = stack.pop() reject_reparse_components(current) directories.append(current) with os.scandir(current) as entries: for entry in entries: if entry.is_symlink() or is_reparse_point(entry.path): raise PrivateFileError(f'private tree contains a link or reparse point: {entry.path}') if entry.is_dir(follow_symlinks=False): stack.append(entry.path) elif entry.is_file(follow_symlinks=False): harden_private_file(entry.path) count += 1 else: raise PrivateFileError(f'unsupported private-tree entry: {entry.path}') for directory in reversed(directories): harden_private_directory(directory) count += 1 return count def require_sensitive_runtime_paths(global_config, create=False): """Verify sensitive runtime parents before any credential or event write.""" config = global_config or {} directories = [] for key in ( 'runtime_dir', 'results_dir', 'result_spool_dir', 'result_bundle_dir', 'queue_dir', 'state_dir', 'log_dir', 'control_dir', 'keycheck_dir', 'postman_cache_dir', 'gharchive_cache_dir', 'work_dir', ): value = config.get(key) if value: directories.append(os.path.abspath(value)) if config.get('postgres_data_dir'): directories.append(canonical_cluster_data_directory({'global': config})) seen = set() for directory in sorted(directories, key=lambda value: (value.count(os.sep), value)): normalized = os.path.normcase(directory) if normalized in seen: continue seen.add(normalized) require_private_directory(directory, create=create) sensitive_files = [config.get('secrets_file')] root_dir = config.get('root_dir') project_dir = config.get('project_dir') if root_dir: sensitive_files.append(os.path.join(root_dir, '.env.postgres')) if project_dir: sensitive_files.append(os.path.join(project_dir, '.env.postgres')) for path in sensitive_files: if not path or not os.path.lexists(path): continue require_protected_sensitive_file_parent( os.path.dirname(os.path.abspath(path)) ) if not private_file_ready(path): raise PrivateFileError(f'sensitive file is not private: {path}') return True def atomic_write_private_json(path, value, max_bytes=MAX_PRIVATE_JSON_BYTES): max_bytes = int(max_bytes) if max_bytes < 1 or max_bytes > MAX_EXTENDED_PRIVATE_JSON_BYTES: raise ValueError('private JSON write bound is invalid') parent = os.path.dirname(os.path.abspath(path)) if parent: require_private_directory(parent, create=True) if os.path.lexists(path): reject_reparse_components(path) payload = json.dumps(value, ensure_ascii=True, sort_keys=True, separators=(',', ':')).encode('utf-8') + b'\n' if len(payload) > max_bytes: raise ValueError('private JSON payload is too large') temporary = f'{path}.{os.getpid()}.{threading.get_ident()}.tmp' flags = os.O_WRONLY | os.O_CREAT | os.O_EXCL if hasattr(os, 'O_BINARY'): flags |= os.O_BINARY descriptor = os.open(temporary, flags, 0o600) try: os.close(descriptor) descriptor = None harden_private_file(temporary) with open(temporary, 'wb') as handle: handle.write(payload) handle.flush() os.fsync(handle.fileno()) if not private_file_ready(temporary): raise PrivateFileError(f'temporary private-file ACL changed: {temporary}') durable_replace(temporary, path) if not private_file_ready(path): raise PrivateFileError(f'private-file ACL changed during replace: {path}') finally: if descriptor is not None: os.close(descriptor) try: if os.path.exists(temporary): os.remove(temporary) except OSError: pass def write_private_json_exclusive(path, value, max_bytes=MAX_PRIVATE_JSON_BYTES): """Atomically publish a private JSON file without replacing an existing name.""" max_bytes = int(max_bytes) if max_bytes < 1 or max_bytes > MAX_EXTENDED_PRIVATE_JSON_BYTES: raise ValueError('private JSON write bound is invalid') parent = os.path.dirname(os.path.abspath(path)) if not parent or not private_directory_ready(parent): raise PrivateFileError(f'private parent directory is not ready: {parent}') reject_reparse_components(parent) payload = json.dumps(value, ensure_ascii=True, sort_keys=True, separators=(',', ':')).encode('utf-8') + b'\n' if len(payload) > max_bytes: raise ValueError('private JSON payload is too large') temporary = f'{path}.{os.getpid()}.{threading.get_ident()}.tmp' flags = os.O_WRONLY | os.O_CREAT | os.O_EXCL if hasattr(os, 'O_BINARY'): flags |= os.O_BINARY descriptor = os.open(temporary, flags, 0o600) try: with os.fdopen(descriptor, 'wb') as handle: descriptor = None handle.write(payload) handle.flush() os.fsync(handle.fileno()) harden_private_file(temporary) if not private_file_ready(temporary): raise PrivateFileError(f'temporary private-file ACL changed: {temporary}') os.link(temporary, path) fsync_directory(parent) if not private_file_ready(path): raise PrivateFileError(f'private-file ACL changed during publication: {path}') finally: if descriptor is not None: os.close(descriptor) try: os.remove(temporary) except OSError: pass def read_private_json(path, max_bytes=MAX_PRIVATE_JSON_BYTES): if not private_file_ready(path): raise PrivateFileError(f'file is absent or not private: {path}') size = os.stat(path, follow_symlinks=False).st_size if size <= 0 or size > max_bytes: raise PrivateFileError(f'private JSON file has invalid size: {path}') with open(path, 'rb') as handle: payload = handle.read(max_bytes + 1) if len(payload) > max_bytes: raise PrivateFileError(f'private JSON file is too large: {path}') try: value = json.loads(payload.decode('utf-8')) except (UnicodeDecodeError, json.JSONDecodeError) as exc: raise PrivateFileError(f'invalid private JSON file: {path}') from exc if not isinstance(value, dict): raise PrivateFileError(f'private JSON root must be an object: {path}') return value