"""Fixed production authority and asynchronous host-operation dispatch.""" import hmac from pathlib import Path, PurePosixPath import threading from host_agent_apply import HostApplyError, HostApplySession from host_agent_lifecycle import execute_fixed_operation from host_agent_protocol import HostAgentStatus, encode_request_payload from host_agent_state import HostOperationState from runtime_document import ( MAX_CONFIG_DOCUMENT_BYTES, _resolve_package_manifest_path, load_yaml_document, ) from runtime_security import read_stable_root_file from scanner_db import ScannerDB from worker_package import ( MAX_WORKER_PACKAGE_MANIFEST_BYTES, load_worker_package_manifest_bytes, ) HOST_RUNTIME_ACTIVE_ROOT = Path('/etc/truf/runtime') HOST_WORKER_PACKAGE_ROOT = Path('/etc/truf/worker-packages') class HostRuntimeError(RuntimeError): def __init__(self, category): self.category = str(category) super().__init__('host operation runtime failed') def _stable_root_file(path, maximum): try: return read_stable_root_file(path, maximum, HOST_WORKER_PACKAGE_ROOT) except Exception: raise HostRuntimeError('package_evidence') from None def _host_manifest_path(resolved): value = PurePosixPath(resolved) try: relative = value.relative_to(PurePosixPath('/data/worker-packages')) except ValueError: raise HostRuntimeError('package_evidence') from None if not relative.parts or any(part in ('', '.', '..') for part in relative.parts): raise HostRuntimeError('package_evidence') return HOST_WORKER_PACKAGE_ROOT.joinpath(*relative.parts) def load_fixed_package_capabilities(config_payload): config = load_yaml_document( config_payload, max_bytes=MAX_CONFIG_DOCUMENT_BYTES, ) try: profiles = config['supervisor']['worker_api']['compatibility_profiles'] except (KeyError, TypeError): raise HostRuntimeError('package_evidence') from None if type(profiles) is not dict: raise HostRuntimeError('package_evidence') evidence = {} try: for profile_name, profile in profiles.items(): if type(profile_name) is not str or type(profile) is not dict: raise HostRuntimeError('package_evidence') reference = profile.get('package_manifest') resolved = _resolve_package_manifest_path(config, reference) if resolved is None: raise HostRuntimeError('package_evidence') payload = _stable_root_file( _host_manifest_path(resolved), MAX_WORKER_PACKAGE_MANIFEST_BYTES, ) manifest = load_worker_package_manifest_bytes(payload) evidence[profile_name] = { 'package_manifest': reference, 'capabilities': manifest['capabilities'], } return evidence except HostRuntimeError: raise except Exception: raise HostRuntimeError('package_evidence') from None finally: config = profiles = profile_name = profile = reference = None resolved = payload = manifest = None class FixedHostOperationDispatcher: def __init__(self): self._guard = threading.Lock() self._active_request = None self._worker = None self._closing = False def _execute(self, session, database, state): try: try: execute_fixed_operation(session, state=state) except BaseException: # The durable executor owns safety/result handling. Do not let # thread tracebacks disclose host details at this outer boundary. pass finally: try: session.close() except BaseException: pass try: database.close() except BaseException: pass finally: with self._guard: self._active_request = None self._worker = None @staticmethod def _record_validation_failure(state, request): try: phase = state.initialize('original') if phase.get('phase') == 'prepared': if phase.get('publication_state') != 'original': return False phase = state.advance( 'prepared', 'failed', 'original', forward_category='validation_failed', safe_detail='validation_failed', ) if ( phase.get('phase') != 'failed' or phase.get('publication_state') != 'original' or phase.get('forward_category') != 'validation_failed' or phase.get('safe_detail') != 'validation_failed' ): return False state.publish_result( 'failed', safe_category='validation_failed', safe_detail='validation_failed', resulting_identity={ 'active_config_sha256': request.active_config_sha256, 'active_secrets_sha256': request.active_secrets_sha256, }, ) return True except Exception: return False def handle(self, request): encoded = encode_request_payload(request) with self._guard: if self._closing: return HostAgentStatus.UNAVAILABLE if self._active_request is not None: return ( HostAgentStatus.ACCEPTED if hmac.compare_digest(encoded, self._active_request) else HostAgentStatus.REJECTED ) state = HostOperationState(request) if state.terminal_result() is not None: return HostAgentStatus.ACCEPTED database = None session = None try: database = ScannerDB.host_agent_authority() session = HostApplySession( request, database, package_capability_provider=load_fixed_package_capabilities, ) session.__enter__() state.initialize(session.publication_state) worker = threading.Thread( target=self._execute, args=(session, database, state), name='truf-host-operation', daemon=False, ) self._active_request = encoded self._worker = worker worker.start() except Exception as error: accepted = ( isinstance(error, HostApplyError) and error.category == 'validation' and session is not None and session.claim is not None and self._record_validation_failure(state, request) ) self._active_request = None self._worker = None if session is not None: try: session.close() except BaseException: pass if database is not None: try: database.close() except BaseException: pass return ( HostAgentStatus.ACCEPTED if accepted else HostAgentStatus.UNAVAILABLE ) return HostAgentStatus.ACCEPTED def close(self): with self._guard: self._closing = True worker = self._worker if worker is not None: worker.join()