215 lines
7.7 KiB
Python
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()
|