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

1613 lines
74 KiB
Python

import copy
from contextlib import contextmanager
import hashlib
import json
import os
from pathlib import Path
import stat
import sys
import tempfile
import traceback
import unittest
from unittest import mock
import yaml
APP_DIR = Path(__file__).resolve().parents[1] / 'app'
if str(APP_DIR) not in sys.path:
sys.path.insert(0, str(APP_DIR))
import runtime_document
import runtime_document_io
import runtime_security
class RuntimeDocumentIOTests(unittest.TestCase):
def documents(self):
template_payload = (APP_DIR / 'config.linux.yaml').read_bytes()
template = runtime_document.load_yaml_document(
template_payload, max_bytes=runtime_document.MAX_CONFIG_DOCUMENT_BYTES,
)
config = copy.deepcopy(template)
docker_pool = config['sources']['dockerhub']['auth_pool']
pools = {
source['auth_pool']
for source in config['sources'].values()
if source.get('auth_pool')
}
secrets = {'auth_pools': {
pool: [{
'name': pool + '_1',
**({'username': 'fixture-user'} if pool == docker_pool else {}),
'token': 'fixture-token',
}]
for pool in pools
}}
return template_payload, config, secrets
def write_documents(self, root, config, secrets):
config_path = root / 'active.yaml'
secrets_path = root / 'secrets.yaml'
config_path.write_text(yaml.safe_dump(config), encoding='utf-8')
secrets_path.write_text(yaml.safe_dump(secrets), encoding='utf-8')
return config_path, secrets_path
@contextmanager
def candidate_store(self, config=None, secrets=None):
_template_payload, default_config, default_secrets = self.documents()
config = copy.deepcopy(default_config if config is None else config)
secrets = copy.deepcopy(default_secrets if secrets is None else secrets)
with tempfile.TemporaryDirectory() as temporary:
root = Path(temporary)
config_path, secrets_path = self.write_documents(root, config, secrets)
candidate_directory = root / 'candidates'
config_candidate = candidate_directory / 'config.yaml'
secrets_candidate = candidate_directory / 'secrets.yaml'
lock_path = candidate_directory / 'candidate.lock'
fsyncs = []
def require_directory(path, create=False):
path = Path(path)
if create:
path.mkdir(mode=0o700, parents=True, exist_ok=True)
if path.is_symlink() or not path.is_dir():
raise OSError('invalid private directory')
return os.fspath(path)
def require_file(path):
path = Path(path)
if path.is_symlink():
raise OSError('invalid private file')
details = path.stat()
if not stat.S_ISREG(details.st_mode):
raise OSError('invalid private file')
return os.fspath(path)
def harden_file(path):
os.chmod(path, 0o600)
class TestLock:
def __init__(self, path):
self.path = Path(path)
def __enter__(self):
self.path.parent.mkdir(mode=0o700, parents=True, exist_ok=True)
self.path.touch(mode=0o600, exist_ok=True)
os.chmod(self.path, 0o600)
return self
def __exit__(self, *_args):
return False
with mock.patch.multiple(
runtime_document_io,
MANAGED_TEMPLATE_PATH=APP_DIR / 'config.linux.yaml',
MANAGED_SECRETS_PATH=secrets_path,
MANAGED_CANDIDATE_DIRECTORY=candidate_directory,
MANAGED_CONFIG_CANDIDATE_PATH=config_candidate,
MANAGED_SECRETS_CANDIDATE_PATH=secrets_candidate,
MANAGED_CANDIDATE_LOCK_PATH=lock_path,
require_private_directory=require_directory,
require_private_file=require_file,
harden_private_file=harden_file,
PrivateFileLock=TestLock,
durable_replace=lambda source, destination: os.replace(source, destination),
fsync_directory=lambda path: fsyncs.append(Path(path)),
):
yield {
'root': root,
'config': config,
'secrets': secrets,
'config_path': config_path,
'secrets_path': secrets_path,
'candidate_directory': candidate_directory,
'config_candidate': config_candidate,
'secrets_candidate': secrets_candidate,
'fsyncs': fsyncs,
}
def preview(self, store, document, payload):
return runtime_document_io.preview_managed_runtime_candidate(
store['config_path'], document, payload,
)
def save(self, store, document, payload, state):
return runtime_document_io.save_managed_runtime_candidate(
store['config_path'],
document,
payload,
expected_active_config_sha256=state.active_config.sha256,
expected_active_secrets_sha256=state.active_secrets.sha256,
expected_candidate_config_sha256=state.candidate_config.sha256,
expected_candidate_secrets_sha256=state.candidate_secrets.sha256,
)
def assert_safe_error(self, error, *sentinels):
rendered = (
str(error),
repr(error),
repr(vars(error)),
repr(error.__cause__),
repr(error.__context__),
''.join(traceback.format_exception(
type(error), error, error.__traceback__,
)),
)
current = error.__traceback__
runtime_locals = []
while current is not None:
if current.tb_frame.f_code.co_filename in {
runtime_document.__file__, runtime_document_io.__file__,
}:
runtime_locals.append(repr(current.tb_frame.f_locals))
current = current.tb_next
for sentinel in sentinels:
text = sentinel.decode('utf-8') if type(sentinel) is bytes else sentinel
self.assertTrue(all(text not in value for value in rendered))
self.assertTrue(all(text not in value for value in runtime_locals))
def test_reads_fixed_private_documents_and_validates_them(self):
_template_payload, config, secrets = self.documents()
with tempfile.TemporaryDirectory() as temporary:
root = Path(temporary)
config_path, secrets_path = self.write_documents(root, config, secrets)
with mock.patch.object(
runtime_document_io, 'require_private_file', side_effect=os.fspath,
) as private_file, mock.patch.object(
runtime_document_io, 'MANAGED_TEMPLATE_PATH', APP_DIR / 'config.linux.yaml',
), mock.patch.object(runtime_document_io, 'MANAGED_SECRETS_PATH', secrets_path):
validated = runtime_document_io.validate_managed_runtime_files(config_path)
self.assertEqual(validated.secrets, secrets)
self.assertIn(mock.call(os.fspath(config_path)), private_file.call_args_list)
self.assertIn(mock.call(os.fspath(APP_DIR / 'config.linux.yaml')), private_file.call_args_list)
self.assertIn(mock.call(os.fspath(secrets_path)), private_file.call_args_list)
self.assertEqual(len(validated.config_sha256), 64)
self.assertEqual(len(validated.secrets_sha256), 64)
def test_manifest_evidence_is_bound_to_the_configured_reference(self):
_template_payload, config, _secrets = self.documents()
worker = config['supervisor']['worker_api']
worker['enabled'] = True
reference = '/data/worker-packages/linux-worker.json'
worker['compatibility_profiles'] = {
'linux': {
'package_manifest': reference,
'sources': ['gitlab', 'dockerhub', 'huggingface'],
},
}
capabilities = [
{'source': 'gitlab', 'platform': 'gitlab', 'planning_kind': 'exact_git_v1'},
{'source': 'dockerhub', 'platform': 'docker', 'planning_kind': 'docker_direct_v1'},
{
'source': 'huggingface', 'platform': 'huggingface',
'planning_kind': 'huggingface_space_v1',
},
]
with mock.patch.object(
runtime_document_io, '_read_package_manifest',
return_value=(b'fixture manifest', None),
) as read_manifest, mock.patch.object(
runtime_document_io, 'load_worker_package_manifest_bytes',
return_value={'capabilities': capabilities},
) as load_manifest:
evidence, ready = runtime_document_io._package_capability_evidence(config)
self.assertTrue(ready)
self.assertEqual(evidence, {
'linux': {'package_manifest': reference, 'capabilities': capabilities},
})
read_manifest.assert_called_once_with(reference)
load_manifest.assert_called_once_with(b'fixture manifest')
config['supervisor']['worker_api']['compatibility_profiles']['linux'][
'package_manifest'
] = '/etc/worker.json'
with mock.patch.object(runtime_document_io, '_read_package_manifest') as read_manifest:
self.assertEqual(
runtime_document_io._package_capability_evidence(config),
(None, None),
)
read_manifest.assert_not_called()
def test_package_manifest_uses_root_owned_reader_on_posix(self):
with mock.patch.object(runtime_document_io.os, 'name', 'posix'), mock.patch.object(
runtime_document_io, 'read_stable_root_file', return_value=b'manifest',
) as stable:
self.assertEqual(
runtime_document_io._read_package_manifest('/data/worker-packages/worker.json'),
(b'manifest', None),
)
stable.assert_called_once_with(
'/data/worker-packages/worker.json',
runtime_document_io.MAX_WORKER_PACKAGE_MANIFEST_BYTES,
runtime_document_io.MANAGED_WORKER_PACKAGE_DIRECTORY,
)
def test_candidate_cannot_redirect_fixed_secrets_read(self):
_template_payload, config, secrets = self.documents()
config['global']['secrets_file'] = '/etc/foreign-secrets.yaml'
with tempfile.TemporaryDirectory() as temporary:
root = Path(temporary)
config_path, secrets_path = self.write_documents(root, config, secrets)
opened = []
def private(path):
opened.append(os.fspath(path))
return os.fspath(path)
with mock.patch.object(runtime_document_io, 'require_private_file', side_effect=private), \
mock.patch.object(
runtime_document_io, 'MANAGED_TEMPLATE_PATH', APP_DIR / 'config.linux.yaml',
), mock.patch.object(
runtime_document_io, 'MANAGED_SECRETS_PATH', secrets_path,
):
with self.assertRaises(runtime_document.RuntimeDocumentError) as raised:
runtime_document_io.validate_managed_runtime_files(config_path)
self.assertEqual(raised.exception.category, 'deployment_path')
self.assertNotIn('/etc/foreign-secrets.yaml', opened)
def test_candidate_error_does_not_retain_secret_bytes(self):
_template_payload, config, _secrets = self.documents()
secret = 'sentinel-io-secret'
with tempfile.TemporaryDirectory() as temporary:
root = Path(temporary)
config_path = root / 'active.yaml'
config_path.write_text(yaml.safe_dump(config), encoding='utf-8')
with mock.patch.object(
runtime_document_io, 'require_private_file', side_effect=os.fspath,
), mock.patch.object(
runtime_document_io, 'MANAGED_TEMPLATE_PATH', APP_DIR / 'config.linux.yaml',
):
with self.assertRaises(runtime_document.RuntimeDocumentError) as raised:
runtime_document_io.validate_managed_runtime_files(
config_path,
secrets_bytes=('auth_pools: [' + secret).encode('utf-8'),
)
error = raised.exception
self.assertNotIn(secret, str(error))
self.assertNotIn(secret, ''.join(traceback.format_exception(
type(error), error, error.__traceback__,
)))
current = error.__traceback__
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document_io.__file__:
self.assertNotIn(secret, repr(current.tb_frame.f_locals))
current = current.tb_next
def test_private_document_rejects_hardlinks_and_changed_identity(self):
with tempfile.TemporaryDirectory() as temporary:
root = Path(temporary)
document = root / 'document.yaml'
document.write_bytes(b'global: {}\n')
linked = root / 'linked.yaml'
try:
os.link(document, linked)
except OSError as exc:
self.skipTest(f'hardlinks unavailable: {exc}')
with mock.patch.object(
runtime_document_io, 'require_private_file', side_effect=os.fspath,
):
payload, failure = runtime_document_io._read_private_document(
linked, 1024, 'config',
)
self.assertIsNone(payload)
self.assertEqual(failure, ('schema', None, None, 'config', 'root'))
linked.unlink()
details = os.stat(document)
current_hardlink = mock.Mock(
st_dev=details.st_dev,
st_ino=details.st_ino,
st_nlink=2,
)
with mock.patch.object(
runtime_document_io, 'require_private_file', side_effect=os.fspath,
), mock.patch.object(
runtime_document_io.os, 'stat', return_value=current_hardlink,
):
payload, failure = runtime_document_io._read_private_document(
document, 1024, 'config',
)
self.assertIsNone(payload)
self.assertEqual(failure, ('schema', None, None, 'config', 'root'))
changed = mock.Mock(st_dev=details.st_dev, st_ino=details.st_ino + 1)
with mock.patch.object(
runtime_document_io, 'require_private_file', side_effect=os.fspath,
), mock.patch.object(runtime_document_io.os, 'stat', return_value=changed):
payload, failure = runtime_document_io._read_private_document(
document, 1024, 'config',
)
self.assertIsNone(payload)
self.assertEqual(failure, ('schema', None, None, 'config', 'root'))
def test_editor_loads_active_then_candidate_without_cross_document_content(self):
with self.candidate_store() as store:
active_config = store['config_path'].read_bytes().decode('utf-8')
active_secrets = store['secrets_path'].read_bytes().decode('utf-8')
config_editor = runtime_document_io.load_managed_runtime_editor_document(
store['config_path'], 'config',
)
secrets_editor = runtime_document_io.load_managed_runtime_editor_document(
store['config_path'], 'secrets',
)
self.assertEqual(config_editor.source, 'active')
self.assertEqual(config_editor.text, active_config)
self.assertNotIn('fixture-token', config_editor.text)
self.assertEqual(secrets_editor.source, 'active')
self.assertEqual(secrets_editor.text, active_secrets)
self.assertNotIn('fixture-token', repr(secrets_editor))
preview = self.preview(store, 'config', active_config.encode('utf-8'))
self.save(store, 'config', active_config.encode('utf-8'), preview.state)
candidate_editor = runtime_document_io.load_managed_runtime_editor_document(
store['config_path'], 'config',
)
self.assertEqual(candidate_editor.source, 'candidate')
self.assertTrue(candidate_editor.selected.present)
def test_candidate_first_save_uses_virtual_active_revision_and_preserves_active(self):
with self.candidate_store() as store:
proposed = copy.deepcopy(store['config'])
proposed['global']['cooldown'] += 1
payload = yaml.safe_dump(proposed, sort_keys=False).encode('utf-8')
active_before = store['config_path'].read_bytes()
preview = self.preview(store, 'config', payload)
self.assertIn(store['candidate_directory'].parent, store['fsyncs'])
self.assertFalse(preview.state.candidate_config.present)
self.assertEqual(
preview.state.candidate_config.sha256,
preview.state.active_config.sha256,
)
revision = self.save(store, 'config', payload, preview.state)
self.assertTrue(revision.created)
self.assertTrue(revision.content_changed)
self.assertTrue(revision.written)
self.assertTrue(revision.after.candidate_config.present)
self.assertEqual(store['config_candidate'].read_bytes(), payload)
self.assertEqual(store['config_path'].read_bytes(), active_before)
self.assertEqual(
revision.proposed.sha256, hashlib.sha256(payload).hexdigest(),
)
if os.name != 'nt':
self.assertEqual(
stat.S_IMODE(store['config_candidate'].stat().st_mode), 0o600,
)
second = self.save(store, 'config', payload, revision.after)
self.assertFalse(second.created)
self.assertFalse(second.content_changed)
self.assertFalse(second.written)
def test_candidate_save_rejects_stale_revision_without_mutation(self):
with self.candidate_store() as store:
payload = yaml.safe_dump(store['config'], sort_keys=False).encode('utf-8')
preview = self.preview(store, 'config', payload)
stale = '0' * 64
with self.assertRaises(runtime_document.RuntimeDocumentError) as raised:
runtime_document_io.save_managed_runtime_candidate(
store['config_path'],
'config',
payload,
expected_active_config_sha256=stale,
expected_active_secrets_sha256=preview.state.active_secrets.sha256,
expected_candidate_config_sha256=preview.state.candidate_config.sha256,
expected_candidate_secrets_sha256=preview.state.candidate_secrets.sha256,
)
self.assertEqual(raised.exception.category, 'reference')
self.assertFalse(store['config_candidate'].exists())
self.assertFalse(store['secrets_candidate'].exists())
def test_invalid_candidates_are_redacted_and_never_staged(self):
cases = []
def unknown_config(store, sentinel):
candidate = copy.deepcopy(store['config'])
candidate['global'][sentinel] = 'hidden-value-' + sentinel
return yaml.safe_dump(candidate, sort_keys=False).encode('utf-8')
def unknown_secrets(store, sentinel):
candidate = copy.deepcopy(store['secrets'])
candidate[sentinel] = {'value': 'hidden-value-' + sentinel}
return yaml.safe_dump(candidate, sort_keys=False).encode('utf-8')
def duplicate_config(_store, sentinel):
return (
f'# {sentinel}\nglobal: {{}}\nglobal: {{}}\n'
).encode('utf-8')
def duplicate_secrets(_store, sentinel):
return (
f'# {sentinel}\nauth_pools: {{}}\nauth_pools: {{}}\n'
).encode('utf-8')
def invalid_reference(store, sentinel):
candidate = copy.deepcopy(store['secrets'])
pool = store['config']['sources']['gitlab']['auth_pool']
candidate['auth_pools'].pop(pool)
return (
yaml.safe_dump(candidate, sort_keys=False) + f'# {sentinel}\n'
).encode('utf-8')
cases.extend((
('unknown-config', 'config', unknown_config, 'unknown_key', 'config'),
('unknown-secrets', 'secrets', unknown_secrets, 'unknown_key', 'secrets'),
('duplicate-config', 'config', duplicate_config, 'duplicate_key', 'config'),
('duplicate-secrets', 'secrets', duplicate_secrets, 'duplicate_key', 'secrets'),
('invalid-reference', 'secrets', invalid_reference, 'reference', 'config'),
))
for name, document, build, category, error_document in cases:
with self.subTest(case=name), self.candidate_store() as store:
baseline = self.preview(
store, 'config', store['config_path'].read_bytes(),
).state
sentinel = f'task55-{name}-sentinel'
payload = build(store, sentinel)
active_config = store['config_path'].read_bytes()
active_secrets = store['secrets_path'].read_bytes()
with self.assertRaises(runtime_document.RuntimeDocumentError) as raised:
self.save(store, document, payload, baseline)
error = raised.exception
self.assertEqual(error.category, category)
self.assertEqual(error.document, error_document)
self.assert_safe_error(error, sentinel, 'hidden-value-' + sentinel)
self.assertEqual(store['config_path'].read_bytes(), active_config)
self.assertEqual(store['secrets_path'].read_bytes(), active_secrets)
self.assertFalse(store['config_candidate'].exists())
self.assertFalse(store['secrets_candidate'].exists())
self.assertEqual(
list(store['candidate_directory'].glob('*.tmp')), [],
)
def test_real_stale_active_revisions_reject_save_and_apply(self):
for changed_document in ('config', 'secrets'):
with self.subTest(document=changed_document), self.candidate_store() as store:
candidate = copy.deepcopy(store['config'])
candidate['global']['cooldown'] += 1
candidate_payload = yaml.safe_dump(
candidate, sort_keys=False,
).encode('utf-8')
preview = self.preview(store, 'config', candidate_payload)
revision = self.save(store, 'config', candidate_payload, preview.state)
stale_state = revision.after
candidate_before = store['config_candidate'].read_bytes()
if changed_document == 'config':
changed = copy.deepcopy(store['config'])
changed['global']['cooldown'] += 2
store['config_path'].write_text(
yaml.safe_dump(changed, sort_keys=False), encoding='utf-8',
)
concurrent_active = store['config_path'].read_bytes()
else:
changed = copy.deepcopy(store['secrets'])
first_pool = next(iter(changed['auth_pools'].values()))
first_pool[0]['token'] = 'concurrent-secret-token'
store['secrets_path'].write_text(
yaml.safe_dump(changed, sort_keys=False), encoding='utf-8',
)
concurrent_active = store['secrets_path'].read_bytes()
with self.assertRaises(runtime_document.RuntimeDocumentError) as save_error:
self.save(store, 'config', candidate_payload, stale_state)
self.assertEqual(save_error.exception.category, 'reference')
self.assertEqual(save_error.exception.path, 'revision')
with self.assertRaises(runtime_document.RuntimeDocumentError) as apply_error:
runtime_document_io.verify_managed_runtime_candidates(
store['config_path'],
'apply-config',
expected_active_config_sha256=stale_state.active_config.sha256,
expected_active_secrets_sha256=stale_state.active_secrets.sha256,
expected_candidate_config_sha256=stale_state.candidate_config.sha256,
)
self.assertEqual(apply_error.exception.category, 'reference')
self.assertEqual(store['config_candidate'].read_bytes(), candidate_before)
active_path = (
store['config_path'] if changed_document == 'config'
else store['secrets_path']
)
self.assertEqual(active_path.read_bytes(), concurrent_active)
def test_stale_existing_candidate_rejects_save_and_apply(self):
with self.candidate_store() as store:
payloads = []
for increment in (1, 2, 3):
candidate = copy.deepcopy(store['config'])
candidate['global']['cooldown'] += increment
payloads.append(yaml.safe_dump(
candidate, sort_keys=False,
).encode('utf-8'))
preview = self.preview(store, 'config', payloads[0])
revision_a = self.save(store, 'config', payloads[0], preview.state)
revision_b = self.save(store, 'config', payloads[1], revision_a.after)
candidate_b = store['config_candidate'].read_bytes()
with self.assertRaises(runtime_document.RuntimeDocumentError) as save_error:
self.save(store, 'config', payloads[2], revision_a.after)
self.assertEqual(save_error.exception.category, 'reference')
with self.assertRaises(runtime_document.RuntimeDocumentError) as apply_error:
runtime_document_io.verify_managed_runtime_candidates(
store['config_path'],
'apply-config',
expected_active_config_sha256=revision_a.after.active_config.sha256,
expected_active_secrets_sha256=revision_a.after.active_secrets.sha256,
expected_candidate_config_sha256=revision_a.after.candidate_config.sha256,
)
self.assertEqual(apply_error.exception.category, 'reference')
self.assertEqual(store['config_candidate'].read_bytes(), candidate_b)
self.assertEqual(
revision_b.after.candidate_config.sha256,
hashlib.sha256(candidate_b).hexdigest(),
)
def test_stale_secrets_candidate_rejects_selected_and_pair_saves(self):
with self.candidate_store() as store:
payloads = []
for token in ('candidate-a', 'candidate-b', 'candidate-c'):
candidate = copy.deepcopy(store['secrets'])
first_pool = next(iter(candidate['auth_pools'].values()))
first_pool[0]['token'] = token
payloads.append(yaml.safe_dump(
candidate, sort_keys=False,
).encode('utf-8'))
preview = self.preview(store, 'secrets', payloads[0])
revision_a = self.save(store, 'secrets', payloads[0], preview.state)
revision_b = self.save(store, 'secrets', payloads[1], revision_a.after)
candidate_b = store['secrets_candidate'].read_bytes()
with self.assertRaises(runtime_document.RuntimeDocumentError):
self.save(store, 'secrets', payloads[2], revision_a.after)
with self.assertRaises(runtime_document.RuntimeDocumentError):
runtime_document_io.verify_managed_runtime_candidates(
store['config_path'],
'apply-secrets',
expected_active_config_sha256=revision_a.after.active_config.sha256,
expected_active_secrets_sha256=revision_a.after.active_secrets.sha256,
expected_candidate_secrets_sha256=revision_a.after.candidate_secrets.sha256,
)
changed_config = copy.deepcopy(store['config'])
changed_config['global']['cooldown'] += 1
with self.assertRaises(runtime_document.RuntimeDocumentError):
self.save(
store,
'config',
yaml.safe_dump(changed_config, sort_keys=False).encode('utf-8'),
revision_a.after,
)
self.assertEqual(store['secrets_candidate'].read_bytes(), candidate_b)
self.assertEqual(
revision_b.after.candidate_secrets.sha256,
hashlib.sha256(candidate_b).hexdigest(),
)
def test_candidate_save_rechecks_state_after_staging_and_cleans_temporary(self):
with self.candidate_store() as store:
candidate = copy.deepcopy(store['config'])
candidate['global']['cooldown'] += 1
payload = yaml.safe_dump(candidate, sort_keys=False).encode('utf-8')
preview = self.preview(store, 'config', payload)
stage = runtime_document_io._stage_candidate
def stage_then_change_active(*args, **kwargs):
temporary = stage(*args, **kwargs)
changed = copy.deepcopy(store['config'])
changed['global']['cooldown'] += 2
store['config_path'].write_text(
yaml.safe_dump(changed, sort_keys=False), encoding='utf-8',
)
return temporary
with mock.patch.object(
runtime_document_io, '_stage_candidate', side_effect=stage_then_change_active,
), mock.patch.object(
runtime_document_io, 'durable_replace', wraps=runtime_document_io.durable_replace,
) as replace, self.assertRaises(runtime_document.RuntimeDocumentError) as raised:
self.save(store, 'config', payload, preview.state)
self.assertEqual(raised.exception.category, 'reference')
replace.assert_not_called()
self.assertFalse(store['config_candidate'].exists())
self.assertEqual(list(store['candidate_directory'].glob('*.tmp')), [])
def test_candidate_save_rolls_back_post_commit_failures_and_races(self):
with self.candidate_store() as store:
candidate = copy.deepcopy(store['config'])
candidate['global']['cooldown'] += 1
payload = yaml.safe_dump(candidate, sort_keys=False).encode('utf-8')
preview = self.preview(store, 'config', payload)
with mock.patch.object(
runtime_document_io,
'_verify_final_candidate',
return_value=runtime_document_io._candidate_failure(
'schema', 'config', 'root',
),
), self.assertRaises(runtime_document.RuntimeDocumentError):
self.save(store, 'config', payload, preview.state)
self.assertFalse(store['config_candidate'].exists())
self.assertEqual(list(store['candidate_directory'].glob('*.tmp')), [])
replace = runtime_document_io.durable_replace
def replace_then_raise(source, destination):
replace(source, destination)
raise OSError('simulated durability failure')
with mock.patch.object(
runtime_document_io,
'durable_replace',
side_effect=replace_then_raise,
), self.assertRaises(runtime_document.RuntimeDocumentError):
self.save(store, 'config', payload, preview.state)
self.assertFalse(store['config_candidate'].exists())
self.assertEqual(list(store['candidate_directory'].glob('*.tmp')), [])
changed = copy.deepcopy(store['config'])
changed['global']['cooldown'] += 2
changed_payload = yaml.safe_dump(changed, sort_keys=False).encode('utf-8')
def replace_then_change_active(source, destination):
replace(source, destination)
store['config_path'].write_bytes(changed_payload)
with mock.patch.object(
runtime_document_io,
'durable_replace',
side_effect=replace_then_change_active,
), self.assertRaises(runtime_document.RuntimeDocumentError) as raised:
self.save(store, 'config', payload, preview.state)
self.assertEqual(raised.exception.category, 'reference')
self.assertEqual(store['config_path'].read_bytes(), changed_payload)
self.assertFalse(store['config_candidate'].exists())
self.assertEqual(list(store['candidate_directory'].glob('*.tmp')), [])
with self.candidate_store() as store:
candidate_a = copy.deepcopy(store['config'])
candidate_a['global']['cooldown'] += 1
payload_a = yaml.safe_dump(candidate_a, sort_keys=False).encode('utf-8')
preview_a = self.preview(store, 'config', payload_a)
revision_a = self.save(store, 'config', payload_a, preview_a.state)
candidate_b = copy.deepcopy(store['config'])
candidate_b['global']['cooldown'] += 2
payload_b = yaml.safe_dump(candidate_b, sort_keys=False).encode('utf-8')
replace = runtime_document_io.durable_replace
def replace_existing_then_raise(source, destination):
replace(source, destination)
raise OSError('simulated durability failure')
with mock.patch.object(
runtime_document_io,
'durable_replace',
side_effect=replace_existing_then_raise,
), self.assertRaises(runtime_document.RuntimeDocumentError):
self.save(store, 'config', payload_b, revision_a.after)
self.assertEqual(store['config_candidate'].read_bytes(), payload_a)
self.assertEqual(list(store['candidate_directory'].glob('*.tmp')), [])
def test_candidate_save_detects_symmetric_post_commit_state_races(self):
for document in ('config', 'secrets'):
with self.subTest(document=document, race='active-counterpart'), self.candidate_store() as store:
if document == 'config':
candidate = copy.deepcopy(store['config'])
candidate['global']['cooldown'] += 1
payload = yaml.safe_dump(candidate, sort_keys=False).encode('utf-8')
concurrent = copy.deepcopy(store['secrets'])
next(iter(concurrent['auth_pools'].values()))[0]['token'] = 'active-race-token'
concurrent_path = store['secrets_path']
else:
candidate = copy.deepcopy(store['secrets'])
next(iter(candidate['auth_pools'].values()))[0]['token'] = 'candidate-token'
payload = yaml.safe_dump(candidate, sort_keys=False).encode('utf-8')
concurrent = copy.deepcopy(store['config'])
concurrent['global']['cooldown'] += 2
concurrent_path = store['config_path']
concurrent_payload = yaml.safe_dump(
concurrent, sort_keys=False,
).encode('utf-8')
preview = self.preview(store, document, payload)
replace = runtime_document_io.durable_replace
def replace_then_change_counterpart(source, destination):
replace(source, destination)
concurrent_path.write_bytes(concurrent_payload)
with mock.patch.object(
runtime_document_io,
'durable_replace',
side_effect=replace_then_change_counterpart,
), self.assertRaises(runtime_document.RuntimeDocumentError) as raised:
self.save(store, document, payload, preview.state)
self.assertEqual(raised.exception.category, 'reference')
self.assertEqual(concurrent_path.read_bytes(), concurrent_payload)
selected = (
store['config_candidate']
if document == 'config'
else store['secrets_candidate']
)
self.assertFalse(selected.exists())
self.assertEqual(list(store['candidate_directory'].glob('*.tmp')), [])
for document in ('config', 'secrets'):
with self.subTest(document=document, race='unselected-candidate'), self.candidate_store() as store:
if document == 'config':
unselected = copy.deepcopy(store['secrets'])
next(iter(unselected['auth_pools'].values()))[0]['token'] = 'unselected-a'
unselected_payload = yaml.safe_dump(
unselected, sort_keys=False,
).encode('utf-8')
unselected_preview = self.preview(store, 'secrets', unselected_payload)
self.save(
store, 'secrets', unselected_payload, unselected_preview.state,
)
proposed = copy.deepcopy(store['config'])
proposed['global']['cooldown'] += 1
concurrent = copy.deepcopy(unselected)
next(iter(concurrent['auth_pools'].values()))[0]['token'] = 'unselected-b'
concurrent_path = store['secrets_candidate']
else:
unselected = copy.deepcopy(store['config'])
unselected['global']['cooldown'] += 1
unselected_payload = yaml.safe_dump(
unselected, sort_keys=False,
).encode('utf-8')
unselected_preview = self.preview(store, 'config', unselected_payload)
self.save(
store, 'config', unselected_payload, unselected_preview.state,
)
proposed = copy.deepcopy(store['secrets'])
next(iter(proposed['auth_pools'].values()))[0]['token'] = 'selected-secrets'
concurrent = copy.deepcopy(unselected)
concurrent['global']['cooldown'] += 1
concurrent_path = store['config_candidate']
proposed_payload = yaml.safe_dump(
proposed, sort_keys=False,
).encode('utf-8')
concurrent_payload = yaml.safe_dump(
concurrent, sort_keys=False,
).encode('utf-8')
preview = self.preview(store, document, proposed_payload)
replace = runtime_document_io.durable_replace
def replace_then_change_unselected(source, destination):
replace(source, destination)
concurrent_path.write_bytes(concurrent_payload)
with mock.patch.object(
runtime_document_io,
'durable_replace',
side_effect=replace_then_change_unselected,
), self.assertRaises(runtime_document.RuntimeDocumentError) as raised:
self.save(store, document, proposed_payload, preview.state)
self.assertEqual(raised.exception.category, 'reference')
self.assertEqual(concurrent_path.read_bytes(), concurrent_payload)
selected = (
store['config_candidate']
if document == 'config'
else store['secrets_candidate']
)
self.assertFalse(selected.exists())
self.assertEqual(list(store['candidate_directory'].glob('*.tmp')), [])
for document in ('config', 'secrets'):
with self.subTest(document=document, race='selected-candidate'), self.candidate_store() as store:
if document == 'config':
proposed = copy.deepcopy(store['config'])
proposed['global']['cooldown'] += 1
concurrent = copy.deepcopy(store['config'])
concurrent['global']['cooldown'] += 2
selected = store['config_candidate']
else:
proposed = copy.deepcopy(store['secrets'])
first = next(iter(proposed['auth_pools'].values()))
first[0]['token'] = 'proposed-token'
concurrent = copy.deepcopy(store['secrets'])
first = next(iter(concurrent['auth_pools'].values()))
first[0]['token'] = 'concurrent-token'
selected = store['secrets_candidate']
proposed_payload = yaml.safe_dump(
proposed, sort_keys=False,
).encode('utf-8')
concurrent_payload = yaml.safe_dump(
concurrent, sort_keys=False,
).encode('utf-8')
preview = self.preview(store, document, proposed_payload)
replace = runtime_document_io.durable_replace
def replace_then_change_selected(source, destination):
replace(source, destination)
selected.write_bytes(concurrent_payload)
with mock.patch.object(
runtime_document_io,
'durable_replace',
side_effect=replace_then_change_selected,
), self.assertRaises(runtime_document.RuntimeDocumentError):
self.save(store, document, proposed_payload, preview.state)
self.assertEqual(selected.read_bytes(), concurrent_payload)
self.assertEqual(list(store['candidate_directory'].glob('*.tmp')), [])
def test_candidate_store_removes_private_crash_temporary_under_lock(self):
with self.candidate_store() as store:
self.preview(store, 'config', store['config_path'].read_bytes())
orphan = store['candidate_directory'] / (
'.secrets.yaml.' + ('a' * 24) + '.tmp'
)
orphan.write_bytes(b'crash-secret-sentinel')
os.chmod(orphan, 0o600)
self.preview(store, 'config', store['config_path'].read_bytes())
self.assertFalse(orphan.exists())
self.assertIn(store['candidate_directory'], store['fsyncs'])
def test_candidate_lock_contention_rejects_without_mutation(self):
with self.candidate_store() as store:
payload = store['config_path'].read_bytes()
preview = self.preview(store, 'config', payload)
active_config = store['config_path'].read_bytes()
active_secrets = store['secrets_path'].read_bytes()
class ContendedLock:
def __init__(self, _path):
pass
def __enter__(self):
raise BlockingIOError('candidate lock is held')
def __exit__(self, *_args):
return False
with mock.patch.object(
runtime_document_io, 'PrivateFileLock', ContendedLock,
), self.assertRaises(runtime_document.RuntimeDocumentError):
self.save(store, 'config', payload, preview.state)
self.assertEqual(store['config_path'].read_bytes(), active_config)
self.assertEqual(store['secrets_path'].read_bytes(), active_secrets)
self.assertFalse(store['config_candidate'].exists())
self.assertFalse(store['secrets_candidate'].exists())
self.assertEqual(list(store['candidate_directory'].glob('*.tmp')), [])
def test_real_candidate_lock_prevents_overlapping_save(self):
with self.candidate_store() as store:
payload = store['config_path'].read_bytes()
preview = self.preview(store, 'config', payload)
runtime_security.harden_private_directory(
os.fspath(store['candidate_directory']),
)
runtime_security.harden_private_file(
os.fspath(store['candidate_directory'] / 'candidate.lock'),
)
with mock.patch.object(
runtime_document_io,
'PrivateFileLock',
runtime_security.PrivateFileLock,
), runtime_security.PrivateFileLock(
store['candidate_directory'] / 'candidate.lock',
), self.assertRaises(runtime_document.RuntimeDocumentError):
self.save(store, 'config', payload, preview.state)
self.assertFalse(store['config_candidate'].exists())
self.assertEqual(list(store['candidate_directory'].glob('*.tmp')), [])
def test_apply_verification_uses_exact_selected_candidate_pair(self):
with self.candidate_store() as store:
proposed_config = copy.deepcopy(store['config'])
proposed_config['global']['cooldown'] += 1
config_payload = yaml.safe_dump(
proposed_config, sort_keys=False,
).encode('utf-8')
config_preview = self.preview(store, 'config', config_payload)
config_revision = self.save(
store, 'config', config_payload, config_preview.state,
)
proposed_secrets = copy.deepcopy(store['secrets'])
first_pool = next(iter(proposed_secrets['auth_pools'].values()))
first_pool[0]['token'] = 'replacement-token'
secrets_payload = yaml.safe_dump(
proposed_secrets, sort_keys=False,
).encode('utf-8')
secrets_preview = self.preview(store, 'secrets', secrets_payload)
secrets_revision = self.save(
store, 'secrets', secrets_payload, secrets_preview.state,
)
state = secrets_revision.after
common = {
'expected_active_config_sha256': state.active_config.sha256,
'expected_active_secrets_sha256': state.active_secrets.sha256,
}
config_only = runtime_document_io.verify_managed_runtime_candidates(
store['config_path'],
'apply-config',
expected_candidate_config_sha256=state.candidate_config.sha256,
**common,
)
self.assertEqual(
(config_only.config_source, config_only.secrets_source),
('candidate', 'active'),
)
secrets_only = runtime_document_io.verify_managed_runtime_candidates(
store['config_path'],
'apply-secrets',
expected_candidate_secrets_sha256=state.candidate_secrets.sha256,
**common,
)
self.assertEqual(
(secrets_only.config_source, secrets_only.secrets_source),
('active', 'candidate'),
)
both = runtime_document_io.verify_managed_runtime_candidates(
store['config_path'],
'apply-both',
expected_candidate_config_sha256=state.candidate_config.sha256,
expected_candidate_secrets_sha256=state.candidate_secrets.sha256,
**common,
)
self.assertEqual(
(both.config_source, both.secrets_source),
('candidate', 'candidate'),
)
self.assertEqual(both.effective_config_sha256, hashlib.sha256(config_payload).hexdigest())
self.assertEqual(both.effective_secrets_sha256, hashlib.sha256(secrets_payload).hexdigest())
self.assertEqual(config_revision.after.candidate_config, state.candidate_config)
def test_apply_verification_rejects_missing_relevant_candidate(self):
with self.candidate_store() as store:
preview = self.preview(
store,
'config',
yaml.safe_dump(store['config']).encode('utf-8'),
)
with self.assertRaises(runtime_document.RuntimeDocumentError) as raised:
runtime_document_io.verify_managed_runtime_candidates(
store['config_path'],
'apply-config',
expected_active_config_sha256=preview.state.active_config.sha256,
expected_active_secrets_sha256=preview.state.active_secrets.sha256,
expected_candidate_config_sha256=preview.state.candidate_config.sha256,
)
self.assertEqual(raised.exception.category, 'reference')
def test_config_diff_is_bounded_and_redacts_sensitive_values(self):
before = {'keychecks': {'env': {'PASSWORD': 'before-secret'}}}
after = {'keychecks': {'env': {'PASSWORD': 'after-secret'}}}
diff = runtime_document_io._config_diff(before, after, 'a', 'b')
self.assertEqual(len(diff.entries), 1)
self.assertTrue(diff.entries[0].value_redacted)
self.assertEqual(diff.entries[0].before, '[redacted]')
self.assertEqual(diff.entries[0].after, '[redacted]')
self.assertNotIn('before-secret', repr(diff))
self.assertNotIn('after-secret', repr(diff))
arbitrary_string = runtime_document_io._config_diff(
{'supervisor': {'defaults': {'extra_args': ['--api-key=raw-secret']}}},
{'supervisor': {'defaults': {'extra_args': ['--api-key=changed-secret']}}},
'a', 'b',
)
self.assertTrue(arbitrary_string.entries[0].value_redacted)
self.assertNotIn('raw-secret', repr(arbitrary_string))
self.assertNotIn('changed-secret', repr(arbitrary_string))
known_sensitive = runtime_document_io._config_diff(
{
'global': {
'database_url': 'postgres://before',
'dashboard_db_url': 'postgres://before-dashboard',
},
'supervisor': {'worker_api': {'admin': {'edge_marker': 'before-marker'}}},
},
{
'global': {
'database_url': 'postgres://after',
'dashboard_db_url': 'postgres://after-dashboard',
},
'supervisor': {'worker_api': {'admin': {'edge_marker': 'after-marker'}}},
},
'a',
'b',
)
self.assertEqual(len(known_sensitive.entries), 3)
self.assertTrue(all(entry.value_redacted for entry in known_sensitive.entries))
self.assertNotIn('postgres://', repr(known_sensitive))
self.assertNotIn('marker', ''.join(
value or ''
for entry in known_sensitive.entries
for value in (entry.before, entry.after)
))
removed_empty = runtime_document_io._config_diff(
{'supervisor': {'defaults': {'extra_args': []}}},
{'supervisor': {'defaults': {}}},
'a',
'b',
)
self.assertEqual(len(removed_empty.entries), 1)
self.assertEqual(removed_empty.entries[0].change, 'removed')
self.assertEqual(removed_empty.entries[0].before, '[]')
self.assertIsNone(removed_empty.entries[0].after)
self.assertFalse(removed_empty.format_only_changed)
many = {f'field_{index}': index for index in range(300)}
bounded = runtime_document_io._config_diff({}, many, 'a', 'b')
self.assertTrue(bounded.truncated)
self.assertLessEqual(
len(bounded.entries), runtime_document_io.MAX_CONFIG_DIFF_ENTRIES,
)
self.assertLessEqual(
bounded.output_bytes,
runtime_document_io.MAX_CONFIG_DIFF_OUTPUT_BYTES,
)
serialized_entries = json.dumps(
[entry.__dict__ for entry in bounded.entries],
ensure_ascii=True,
sort_keys=True,
separators=(',', ':'),
).encode('utf-8')
self.assertEqual(len(serialized_entries), bounded.output_bytes)
def test_secrets_preview_returns_aggregate_only(self):
pool_sentinel = 'pool-sentinel'
entry_sentinel = 'entry-sentinel'
username_sentinel = 'username-sentinel'
token_sentinel = 'token-sentinel'
with self.candidate_store() as store:
proposed = copy.deepcopy(store['secrets'])
proposed['auth_pools'][pool_sentinel] = [{
'name': entry_sentinel,
'username': username_sentinel,
'token': token_sentinel,
}]
payload = yaml.safe_dump(proposed, sort_keys=False).encode('utf-8')
preview = self.preview(store, 'secrets', payload)
rendered = repr(preview.diff)
self.assertEqual(preview.diff.pools_added, 1)
self.assertEqual(preview.diff.entries_added, 1)
for sentinel in (
pool_sentinel, entry_sentinel, username_sentinel, token_sentinel,
):
self.assertNotIn(sentinel, rendered)
def test_candidate_symlink_and_hardlink_are_rejected(self):
with self.candidate_store() as store:
store['candidate_directory'].mkdir(parents=True)
foreign = store['root'] / 'foreign.yaml'
foreign.write_bytes(store['config_path'].read_bytes())
try:
store['config_candidate'].symlink_to(foreign)
except OSError as exc:
self.skipTest(f'symlinks unavailable: {exc}')
with self.assertRaises(runtime_document.RuntimeDocumentError):
self.preview(
store, 'config', store['config_path'].read_bytes(),
)
store['config_candidate'].unlink()
try:
os.link(foreign, store['config_candidate'])
except OSError as exc:
self.skipTest(f'hardlinks unavailable: {exc}')
with self.assertRaises(runtime_document.RuntimeDocumentError):
self.preview(
store, 'config', store['config_path'].read_bytes(),
)
def test_candidate_failure_traceback_does_not_retain_secret_bytes(self):
secret = 'candidate-traceback-sentinel'
with self.candidate_store() as store:
payload = ('auth_pools: [' + secret).encode('utf-8')
with self.assertRaises(runtime_document.RuntimeDocumentError) as raised:
self.preview(store, 'secrets', payload)
error = raised.exception
self.assertNotIn(secret, str(error))
self.assertNotIn(secret, ''.join(traceback.format_exception(
type(error), error, error.__traceback__,
)))
current = error.__traceback__
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document_io.__file__:
self.assertNotIn(secret, repr(current.tb_frame.f_locals))
current = current.tb_next
def test_cancellation_tracebacks_do_not_retain_candidate_bytes(self):
secret = b'cancellation-candidate-sentinel'
with mock.patch.object(
runtime_document_io, '_read_private_document',
side_effect=KeyboardInterrupt,
):
try:
runtime_document_io._validate_payload_pair(secret, secret)
except KeyboardInterrupt as error:
current = error.__traceback__
else:
self.fail('KeyboardInterrupt was not propagated')
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document_io.__file__:
self.assertNotIn(repr(secret), repr(current.tb_frame.f_locals))
current = current.tb_next
_template_payload, config, secrets = self.documents()
template = runtime_document.load_yaml_document(
_template_payload,
max_bytes=runtime_document.MAX_CONFIG_DOCUMENT_BYTES,
)
config_secret = 'config-cancellation-edge-marker'
config['supervisor']['worker_api']['admin']['edge_marker'] = config_secret
original_bounded_text = runtime_document._bounded_text
def interrupt_on_config_secret(value, *args, **kwargs):
if value == config_secret:
raise KeyboardInterrupt
return original_bounded_text(value, *args, **kwargs)
with mock.patch.object(
runtime_document,
'_bounded_text',
side_effect=interrupt_on_config_secret,
):
try:
runtime_document.validate_runtime_documents(
config,
secrets,
config_template=template,
package_capabilities={},
)
except KeyboardInterrupt as error:
current = error.__traceback__
else:
self.fail('KeyboardInterrupt was not propagated')
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document.__file__:
self.assertNotIn(config_secret, repr(current.tb_frame.f_locals))
current = current.tb_next
_template_payload, config, secrets = self.documents()
template = runtime_document.load_yaml_document(
_template_payload,
max_bytes=runtime_document.MAX_CONFIG_DOCUMENT_BYTES,
)
env_secret = 'config-cancellation-environment-value'
config['supervisor']['sources']['gitlab']['env'] = {
'FIXTURE_CREDENTIAL': env_secret,
}
def interrupt_on_env_secret(value, *args, **kwargs):
if value == env_secret:
raise KeyboardInterrupt
return original_bounded_text(value, *args, **kwargs)
with mock.patch.object(
runtime_document,
'_bounded_text',
side_effect=interrupt_on_env_secret,
):
try:
runtime_document.validate_runtime_documents(
config,
secrets,
config_template=template,
package_capabilities={},
)
except KeyboardInterrupt as error:
current = error.__traceback__
else:
self.fail('KeyboardInterrupt was not propagated')
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document.__file__:
self.assertNotIn(env_secret, repr(current.tb_frame.f_locals))
current = current.tb_next
evidence_config = {
'supervisor': {'worker_api': {'compatibility_profiles': {
'fixture': {'package_manifest': '/data/worker-packages/worker-package.json'},
}}},
'sensitive': config_secret,
}
with mock.patch.object(
runtime_document_io,
'_resolve_package_manifest_path',
side_effect=KeyboardInterrupt,
):
try:
runtime_document_io._package_capability_evidence(evidence_config)
except KeyboardInterrupt as error:
current = error.__traceback__
else:
self.fail('KeyboardInterrupt was not propagated')
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document_io.__file__:
self.assertNotIn(config_secret, repr(current.tb_frame.f_locals))
current = current.tb_next
_template_payload, config, _secrets = self.documents()
config['supervisor']['worker_api']['admin']['edge_marker'] = config_secret
semantic_interruptions = (
(
'_validate_required_fields',
mock.patch.object(
runtime_document,
'_get_path',
side_effect=KeyboardInterrupt,
),
lambda: runtime_document._validate_required_fields(config),
),
(
'_validate_worker_config',
mock.patch.object(
runtime_document,
'_normalize_capabilities',
side_effect=KeyboardInterrupt,
),
lambda: runtime_document._validate_worker_config(config, {}),
),
(
'_validate_deployment_paths',
mock.patch.object(
runtime_document,
'_expand_path',
side_effect=KeyboardInterrupt,
),
lambda: runtime_document._validate_deployment_paths(config),
),
)
for name, interruption, call in semantic_interruptions:
with self.subTest(config_cancellation=name), interruption:
try:
call()
except KeyboardInterrupt as error:
current = error.__traceback__
else:
self.fail('KeyboardInterrupt was not propagated')
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document.__file__:
self.assertNotIn(config_secret, repr(current.tb_frame.f_locals))
current = current.tb_next
class InterruptingPath:
def __iter__(self):
raise KeyboardInterrupt
try:
runtime_document._get_path(config, InterruptingPath())
except KeyboardInterrupt as error:
current = error.__traceback__
else:
self.fail('KeyboardInterrupt was not propagated')
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document.__file__:
self.assertNotIn(config_secret, repr(current.tb_frame.f_locals))
current = current.tb_next
loader_secret = 'yaml-loader-cancellation-secret'
loader_payload = (
'auth_pools:\n fixture:\n - name: fixture\n'
f' token: {loader_secret}\n'
).encode('utf-8')
construct_object = runtime_document._StrictSafeLoader.construct_object
def interrupt_loader(self, node, deep=False):
if getattr(node, 'value', None) == loader_secret:
raise KeyboardInterrupt
return construct_object(self, node, deep=deep)
with mock.patch.object(
runtime_document._StrictSafeLoader,
'construct_object',
new=interrupt_loader,
):
try:
runtime_document.load_yaml_document(
loader_payload,
max_bytes=runtime_document.MAX_SECRETS_DOCUMENT_BYTES,
)
except KeyboardInterrupt as error:
current = error.__traceback__
else:
self.fail('KeyboardInterrupt was not propagated')
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document.__file__:
self.assertNotIn(loader_secret, repr(current.tb_frame.f_locals))
current = current.tb_next
schema_config = copy.deepcopy(config)
schema_config['supervisor']['worker_api']['admin']['edge_marker'] = config_secret
dynamic_schema_kind = runtime_document._dynamic_schema_kind
def interrupt_schema_identity(path):
if path[-1:] == ('edge_marker',):
raise KeyboardInterrupt
return dynamic_schema_kind(path)
with mock.patch.object(
runtime_document,
'_dynamic_schema_kind',
side_effect=interrupt_schema_identity,
):
try:
runtime_document._config_template_schema_hash(schema_config)
except KeyboardInterrupt as error:
current = error.__traceback__
else:
self.fail('KeyboardInterrupt was not propagated')
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document.__file__:
self.assertNotIn(config_secret, repr(current.tb_frame.f_locals))
current = current.tb_next
preliminary_payload = yaml.safe_dump(
schema_config, sort_keys=False,
).encode('utf-8')
with mock.patch.object(
runtime_document_io,
'_read_private_document',
return_value=(preliminary_payload, None),
), mock.patch.object(
runtime_document_io.hashlib,
'sha256',
side_effect=KeyboardInterrupt,
):
try:
runtime_document_io._load_managed_runtime_config_inner(
Path('active.yaml'),
)
except KeyboardInterrupt as error:
current = error.__traceback__
else:
self.fail('KeyboardInterrupt was not propagated')
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document_io.__file__:
self.assertNotIn(config_secret, repr(current.tb_frame.f_locals))
current = current.tb_next
_template_payload, config, secrets = self.documents()
decoded_secret = 'decoded-cancellation-secret'
first_pool = next(iter(secrets['auth_pools'].values()))
first_pool[0]['token'] = decoded_secret
original_bounded_text = runtime_document._bounded_text
def interrupt_on_decoded_secret(value, *args, **kwargs):
if value == decoded_secret:
raise KeyboardInterrupt
return original_bounded_text(value, *args, **kwargs)
with mock.patch.object(
runtime_document,
'_bounded_text',
side_effect=interrupt_on_decoded_secret,
):
try:
runtime_document.validate_runtime_documents(
config,
secrets,
config_template=runtime_document.load_yaml_document(
_template_payload,
max_bytes=runtime_document.MAX_CONFIG_DOCUMENT_BYTES,
),
package_capabilities={},
)
except KeyboardInterrupt as error:
current = error.__traceback__
else:
self.fail('KeyboardInterrupt was not propagated')
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document.__file__:
self.assertNotIn(decoded_secret, repr(current.tb_frame.f_locals))
current = current.tb_next
template = runtime_document.load_yaml_document(
_template_payload,
max_bytes=runtime_document.MAX_CONFIG_DOCUMENT_BYTES,
)
for target in ('_consume', '_validate_worker_config'):
with self.subTest(decoded_cancellation=target), mock.patch.object(
runtime_document,
target,
side_effect=KeyboardInterrupt,
):
try:
if target == '_consume':
runtime_document._normalize_secrets(secrets)
else:
runtime_document.validate_runtime_documents(
config,
secrets,
config_template=template,
package_capabilities={},
)
except KeyboardInterrupt as error:
current = error.__traceback__
else:
self.fail('KeyboardInterrupt was not propagated')
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document.__file__:
self.assertNotIn(
decoded_secret, repr(current.tb_frame.f_locals),
)
current = current.tb_next
with self.candidate_store() as store:
store['candidate_directory'].mkdir(parents=True, exist_ok=True)
real_sha256 = hashlib.sha256
hash_calls = 0
def interrupt_after_staging(payload):
nonlocal hash_calls
hash_calls += 1
if hash_calls == 2:
raise KeyboardInterrupt
return real_sha256(payload)
with mock.patch.object(
runtime_document_io.hashlib,
'sha256',
side_effect=interrupt_after_staging,
):
try:
runtime_document_io._stage_candidate(
store['config_candidate'], secret, len(secret),
)
except KeyboardInterrupt as error:
current = error.__traceback__
else:
self.fail('KeyboardInterrupt was not propagated')
self.assertFalse(store['config_candidate'].exists())
self.assertEqual(list(store['candidate_directory'].glob('*.tmp')), [])
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document_io.__file__:
self.assertNotIn(repr(secret), repr(current.tb_frame.f_locals))
current = current.tb_next
validated = runtime_document.ValidatedRuntimeDocuments(
config={'ordinary': secret.decode('ascii')},
secrets={'auth_pools': {}},
)
interrupted_helpers = (
(
'_document_identity',
mock.patch.object(
runtime_document_io.hashlib, 'sha256',
side_effect=KeyboardInterrupt,
),
lambda: runtime_document_io._document_identity(secret),
),
(
'_config_diff',
mock.patch.object(
runtime_document_io, '_render_diff_value',
side_effect=KeyboardInterrupt,
),
lambda: runtime_document_io._config_diff(
{'ordinary': 1, 'secret_holder': secret.decode('ascii')},
{'ordinary': 2, 'secret_holder': secret.decode('ascii')},
'a',
'b',
),
),
(
'_candidate_diff',
mock.patch.object(
runtime_document_io, '_secrets_diff',
side_effect=KeyboardInterrupt,
),
lambda: runtime_document_io._candidate_diff(
'secrets', validated, validated, secret, secret,
),
),
)
for name, interruption, call in interrupted_helpers:
with self.subTest(helper=name), interruption:
try:
call()
except KeyboardInterrupt as error:
current = error.__traceback__
else:
self.fail('KeyboardInterrupt was not propagated')
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document_io.__file__:
self.assertNotIn(repr(secret), repr(current.tb_frame.f_locals))
current = current.tb_next
public_calls = (
(
runtime_document_io,
'_validate_managed_runtime_files_inner',
lambda: runtime_document_io.validate_managed_runtime_files(
Path('active.yaml'), secrets_bytes=secret,
),
),
(
runtime_document_io,
'_preview_candidate_inner',
lambda: runtime_document_io.preview_managed_runtime_candidate(
Path('active.yaml'), 'secrets', secret,
),
),
(
runtime_document_io,
'_save_candidate_inner',
lambda: runtime_document_io.save_managed_runtime_candidate(
Path('active.yaml'),
'secrets',
secret,
expected_active_config_sha256='0' * 64,
expected_active_secrets_sha256='0' * 64,
expected_candidate_config_sha256='0' * 64,
expected_candidate_secrets_sha256='0' * 64,
),
),
(
runtime_document,
'_preview_runtime_documents_inner',
lambda: runtime_document.preview_runtime_documents(
secret,
secret,
config_template_payload=secret,
),
),
)
for module, target, call in public_calls:
with self.subTest(target=target), mock.patch.object(
module, target, side_effect=KeyboardInterrupt,
):
try:
call()
except KeyboardInterrupt as error:
current = error.__traceback__
else:
self.fail('KeyboardInterrupt was not propagated')
while current is not None:
if current.tb_frame.f_code.co_filename in {
runtime_document_io.__file__, runtime_document.__file__,
}:
self.assertNotIn(repr(secret), repr(current.tb_frame.f_locals))
current = current.tb_next
with mock.patch.object(runtime_document_io.os, 'open', side_effect=KeyboardInterrupt):
try:
runtime_document_io._stage_candidate(
Path('candidate.yaml'), secret, len(secret),
)
except KeyboardInterrupt as error:
current = error.__traceback__
else:
self.fail('KeyboardInterrupt was not propagated')
while current is not None:
if current.tb_frame.f_code.co_filename == runtime_document_io.__file__:
self.assertNotIn(repr(secret), repr(current.tb_frame.f_locals))
current = current.tb_next
if __name__ == '__main__':
unittest.main()