353 lines
16 KiB
Python
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()
|