Files
SystemSimulationApp/tests/test_result_storage.py
T

353 lines
16 KiB
Python

"""Real filesystem/lease tests for result quota, eviction and crash inspection."""
from concurrent.futures import ThreadPoolExecutor
import json
from itertools import count
import os
from pathlib import Path
import shutil
import struct
import subprocess
import sys
import tempfile
import threading
from types import SimpleNamespace
import unittest
from unittest.mock import patch
import zlib
from app.simulation.native_codegen import result_storage as storage
def block(sequence=0, times=(0., .1), values=(1., -0.)):
payload = struct.pack('<' + 'd' * (len(times) + len(values)), *times, *values)
return struct.pack('<8sQQQQ', b'SIMBLK01', sequence, len(times), 2, zlib.crc32(payload)) + payload + b'COMMIT01'
class ResultStorageTests(unittest.TestCase):
def setUp(self):
self.temp = tempfile.TemporaryDirectory(prefix='result storage ')
self.addCleanup(self.temp.cleanup)
self.root = Path(self.temp.name) / 'simresults'
# Distinct lease shards make eviction order assertions deterministic.
identifiers = count(1)
patcher = patch.object(storage, 'uuid4', side_effect=lambda:
SimpleNamespace(hex=f'{next(identifiers):08x}' + '0' * 24))
patcher.start()
self.addCleanup(patcher.stop)
def archive(self, limit=300000):
archive = storage.ResultArchive({'stateKeys': ['x'], 'variables': [{'key': 'x'}]},
root=self.root, limit_bytes=limit)
self.addCleanup(archive.close)
return archive
def test_default_and_configuration(self):
with patch.dict(os.environ, {}, clear=True):
self.assertEqual(storage.result_limit_bytes(), 1024**3)
with patch.dict(os.environ, {'SIMULATION_RESULT_STORAGE_MB': '2'}):
self.assertEqual(storage.result_limit_bytes(), 2 * 1024**2)
for value in ('0', '-1', 'abc', '1.5'):
with patch.dict(os.environ, {'SIMULATION_RESULT_STORAGE_MB': value}):
with self.assertRaises(ValueError):
storage.result_limit_bytes()
def test_archive_reference_is_relative_and_legacy_absolute_hint_is_ignored(self):
archive = self.archive()
summary = archive.finish()
self.assertEqual(summary['directory'], f'simresults/{archive.path.name}')
self.assertEqual(summary['directoryBase'], 'project')
self.assertNotIn(str(self.root), (archive.path / 'manifest.json').read_text(encoding='utf-8'))
data = json.loads((archive.path / 'manifest.json').read_bytes())
data['blocks']['directory'] = 'Z:\\another-machine\\old-project\\simresults\\' + archive.path.name
(archive.path / 'manifest.json').write_text(json.dumps(data), encoding='utf-8')
self.assertEqual(storage.result_archive_path(data['id'], root=self.root), archive.path)
for invalid in ('../outside', '/outside', 'C:\\outside', summary['directory']):
with self.assertRaises(ValueError):
storage.result_archive_path(invalid, root=self.root)
def test_moved_project_imported_from_other_cwd_finds_and_writes_results(self):
base = Path(self.temp.name)
original = base / 'original project'
moved = base / '移动后的 项目'
package = original / 'app' / 'simulation' / 'native_codegen'
package.mkdir(parents=True)
for directory in (package, package.parent, package.parent.parent):
(directory / '__init__.py').write_text('', encoding='utf-8')
source = Path(storage.__file__).parent
for name in ('result_storage.py', 'cache_storage.py'):
shutil.copyfile(source / name, package / name)
archive = storage.ResultArchive({'stateKeys': ['x']}, root=original / 'simresults')
try:
archive.reserve(len(block()))
(archive.path / 'states.bin').write_bytes(block())
summary = archive.finish()
finally:
archive.close()
# Only move this test's own directory within its verified temporary root.
self.assertEqual(original.resolve().parent, base.resolve())
self.assertEqual(moved.resolve().parent, base.resolve())
original.rename(moved)
elsewhere = base / 'unrelated cwd'
elsewhere.mkdir()
script = '''import sys,json
from pathlib import Path
sys.path.insert(0,sys.argv[1])
from app.simulation.native_codegen import result_storage as s
project=Path(sys.argv[1])
assert s.RESULT_ROOT == project/'simresults'
old=s.result_archive_path(sys.argv[2])
assert s.scan_blocks(old/'states.bin')['samples']==2
assert json.loads((old/'manifest.json').read_bytes())['blocks']['directory']==f'simresults/{sys.argv[2]}'
new=s.ResultArchive({})
try:
assert new.path.parent==project/'simresults'
assert new.finish()['directory']==f'simresults/{new.path.name}'
finally:
new.close()
assert not (Path.cwd()/'simresults').exists()
print('relocation-ok')
'''
process = subprocess.run([sys.executable, '-c', script, str(moved), summary['id']],
cwd=elsewhere, capture_output=True, text=True, timeout=30)
self.assertEqual(process.returncode, 0, process.stderr)
self.assertEqual(process.stdout.strip(), 'relocation-ok')
self.assertFalse(original.exists())
def test_oldest_finished_archive_is_deleted_and_active_run_is_pinned(self):
old = self.archive()
old.reserve(70000)
(old.path / 'states.bin').write_bytes(b'x' * 70000)
old.finish()
newer = self.archive()
newer.reserve(70000)
(newer.path / 'states.bin').write_bytes(b'x' * 70000)
newer.finish()
current = self.archive()
current.reserve(110000)
self.assertFalse(old.path.exists())
self.assertTrue(newer.path.exists())
self.assertTrue(current.path.exists())
with self.assertRaises(storage.ResultQuotaError):
current.reserve(300000)
self.assertTrue(current.path.exists())
def test_concurrent_reservations_cannot_overcommit(self):
first, second = self.archive(250000), self.archive(250000)
barrier = threading.Barrier(2)
def reserve(archive):
barrier.wait()
try:
archive.reserve(100000)
return True
except storage.ResultQuotaError:
return False
with ThreadPoolExecutor(2) as pool:
results = list(pool.map(reserve, (first, second)))
self.assertEqual(sorted(results), [False, True])
self.assertTrue(first.path.exists() and second.path.exists())
def test_reserved_window_has_no_per_block_filesystem_work(self):
archive = self.archive(storage.DEFAULT_LIMIT_BYTES)
archive.reserve(4096)
persisted = (archive.path / 'manifest.json').read_bytes()
granted = archive.data['reservedBytes'] - archive.metadata_budget
self.assertEqual(granted, 16 * 1024**2)
with patch.object(archive, '_locked', side_effect=AssertionError('unexpected quota lock')), \
patch.object(archive, '_make_room', side_effect=AssertionError('unexpected scan')), \
patch.object(archive, '_publish', side_effect=AssertionError('unexpected write')):
for _ in range(999):
archive.reserve(4096)
self.assertEqual(archive._credit, granted - 1000 * 4096)
self.assertEqual((archive.path / 'manifest.json').read_bytes(), persisted)
def test_refill_reserves_only_missing_bytes_and_publishes_new_window(self):
archive = self.archive(8 * 1024**2)
window = archive._reservation_window
archive.reserve(window - 20)
with patch.object(archive, '_make_room', wraps=archive._make_room) as scan, \
patch.object(archive, '_publish', wraps=archive._publish) as publish:
archive.reserve(45)
scan.assert_called_once_with(25)
publish.assert_called_once_with()
self.assertEqual(archive._credit, window - 25)
manifest = json.loads((archive.path / 'manifest.json').read_bytes())
self.assertEqual(manifest['reservedBytes'], archive.metadata_budget + 2 * window)
def test_optional_credit_does_not_evict_history_or_reject_fitting_block(self):
old = self.archive()
old.finish()
current = self.archive()
with current._locked():
free = current._make_room(1)
unknown = self.root / 'user.bin'
unknown.write_bytes(b'x' * (free - 100))
current.reserve(80)
self.assertEqual(current._credit, 20)
self.assertTrue(old.path.exists())
current.reserve(20)
self.assertTrue(old.path.exists())
self.assertEqual(unknown.stat().st_size, free - 100)
def test_failed_publication_grants_no_credit_and_can_retry(self):
archive = self.archive()
archive.reserve(100)
previous = archive.data['reservedBytes']
credit = archive._credit
persisted = (archive.path / 'manifest.json').read_bytes()
with patch.object(archive, '_publish', side_effect=OSError('disk failure')):
with self.assertRaises(OSError):
archive.reserve(credit + 1)
self.assertEqual(archive.data['reservedBytes'], previous)
self.assertEqual(archive._credit, credit)
self.assertEqual((archive.path / 'manifest.json').read_bytes(), persisted)
archive.reserve(credit + 1)
self.assertEqual(archive._credit, archive._reservation_window - 1)
def test_threads_share_credit_without_lost_deductions(self):
archive = self.archive(64 * 1024**2)
with ThreadPoolExecutor(8) as pool:
list(pool.map(archive.reserve, [8192] * 1000))
granted = archive.data['reservedBytes'] - archive.metadata_budget
self.assertEqual(granted - archive._credit, 8192 * 1000)
self.assertGreaterEqual(archive._credit, 0)
def test_closed_archive_cannot_spend_leftover_credit(self):
archive = self.archive()
archive.reserve(1)
self.assertGreater(archive._credit, 0)
archive.finish()
manifest = json.loads((archive.path / 'manifest.json').read_bytes())
self.assertEqual(manifest['reservedBytes'], archive.metadata_budget + 1)
with self.assertRaises(OSError):
archive.reserve(1)
def test_failed_refill_preserves_credit_for_a_smaller_final_block(self):
archive = self.archive()
with archive._locked():
free = archive._make_room(1)
(self.root / 'user.bin').write_bytes(b'x' * (free - 100))
archive.reserve(80)
with self.assertRaises(storage.ResultQuotaError):
archive.reserve(21)
self.assertEqual(archive._credit, 20)
archive.reserve(20)
self.assertEqual(archive._credit, 0)
with self.assertRaises(storage.ResultQuotaError):
archive.reserve(1)
def test_unused_credit_is_reusable_after_finish(self):
first = self.archive(250000)
first.reserve(1)
first.finish()
second = self.archive(250000)
with second._locked():
free = second._make_room(1)
second.reserve(free)
self.assertTrue(first.path.exists())
self.assertEqual(second._credit, 0)
with self.assertRaises(storage.ResultQuotaError):
# A user addition makes the old archive non-evictable.
(first.path / 'keep.txt').write_bytes(b'keep')
second.reserve(1)
def test_another_process_accounts_for_unused_credit(self):
first = self.archive(250000)
first.reserve(1)
# This fits if only consumed bytes were counted, but not when the
# persisted, unspent quota window is correctly pinned by our lease.
amount = 250000 - 2 * first.metadata_budget - 1000
script = '''import sys
from pathlib import Path
from app.simulation.native_codegen.result_storage import ResultArchive, ResultQuotaError
archive = ResultArchive({}, root=Path(sys.argv[1]), limit_bytes=250000)
try:
try:
archive.reserve(int(sys.argv[2]))
except ResultQuotaError:
print('quota-protected')
else:
raise AssertionError('unspent credit was reused')
finally:
archive.close()
'''
process = subprocess.run([sys.executable, '-c', script, str(self.root), str(amount)],
capture_output=True, text=True, timeout=30)
self.assertEqual(process.returncode, 0, process.stderr)
self.assertEqual(process.stdout.strip(), 'quota-protected')
def test_other_process_lease_blocks_eviction(self):
archive = self.archive()
path, key = archive.path, archive.key
archive.finish()
code = ('import sys; from pathlib import Path; '
'from app.simulation.native_codegen.cache_storage import acquire_cache_lease; '
'lease=acquire_cache_lease(Path(sys.argv[1]),"models",sys.argv[2]); '
'print("ready",flush=True); sys.stdin.readline(); lease.close()')
child = subprocess.Popen([sys.executable, '-c', code, str(self.root), key],
stdin=subprocess.PIPE, stdout=subprocess.PIPE, text=True)
try:
self.assertEqual(child.stdout.readline().strip(), 'ready')
current = self.archive(140000)
with self.assertRaises(storage.ResultQuotaError):
current.reserve(10000)
self.assertTrue(path.exists())
finally:
child.communicate('\n', timeout=10)
def test_unknown_content_is_counted_and_never_deleted(self):
self.root.mkdir()
unknown = self.root / 'important.txt'
unknown.write_bytes(b'a' * 100000)
with self.assertRaises(storage.ResultQuotaError):
self.archive(120000)
self.assertEqual(unknown.stat().st_size, 100000)
def test_archive_with_user_added_files_is_not_deleted(self):
old = self.archive()
old.finish()
added = old.path / 'user.txt'
added.write_bytes(b'a' * 100000)
with self.assertRaises(storage.ResultQuotaError):
self.archive(120000)
self.assertTrue(added.exists())
def test_committed_prefix_survives_truncation_and_detects_corruption(self):
archive = self.archive()
path = archive.path / 'states.bin'
first, second = block(), block(1, (.2, .3))
archive.reserve(len(first) + len(second))
path.write_bytes(first + second[:-3])
result = storage.scan_blocks(path, expected_columns=2)
self.assertEqual(result['samples'], 2)
self.assertEqual(result['storedUntil'], .1)
self.assertEqual(result['validBytes'], len(first))
self.assertTrue(result['incompleteTail'])
damaged = bytearray(second)
damaged[50] ^= 1
path.write_bytes(first + damaged)
result = storage.scan_blocks(path, expected_columns=2)
self.assertEqual(result['samples'], 2)
self.assertTrue(result['corrupt'])
summary = archive.finish()
self.assertEqual(summary['status'], 'failed')
manifest = json.loads((archive.path / 'manifest.json').read_bytes())
self.assertEqual(manifest['blocks']['states']['samples'], 2)
def test_orphan_worker_is_not_evicted_until_process_exit(self):
old = self.archive()
child = subprocess.Popen([sys.executable, '-c', 'import sys;sys.stdin.readline()'], stdin=subprocess.PIPE)
try:
old.worker_started(child.pid)
# Simulate loss of the parent's lease while its worker still runs.
old.lease.close()
current = self.archive(140000)
with self.assertRaises(storage.ResultQuotaError):
current.reserve(20000)
self.assertTrue(old.path.exists())
finally:
child.communicate(b'\n', timeout=10)
if __name__ == '__main__':
unittest.main()