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

215 lines
7.7 KiB
Python

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