"""Small synthetic checks only; model/real-source runs belong on Worker 006.""" import importlib import io import json import threading import zipfile from pathlib import Path import numpy as np import pytest @pytest.fixture def pilot(monkeypatch): root = Path(__file__).resolve().parents[1] monkeypatch.syspath_prepend( str(root / "experiments/perception/worker/streaming_profile_stage1") ) return lambda name: importlib.import_module(name) def test_pending_overflow_is_explicit_and_does_not_evict_active(pilot): queue = pilot("pilot_queue").Mailbox(capacity=2, byte_limit=100) def make(sequence): return {"sequence": sequence, "payload_bytes": 10} queue.put(make(0)) active = queue.take() for sequence in (1, 2, 3): queue.put(make(sequence)) assert queue.dropped == [{"sequence": 1, "reason": "pending-overflow"}] assert queue.bytes == 30 assert queue.peak_pending == 2 queue.release(active) assert queue.take()["sequence"] == 2 assert queue.take()["sequence"] == 3 def test_active_bytes_count_towards_memory_limit(pilot): queue = pilot("pilot_queue").Mailbox(byte_limit=20) queue.put({"sequence": 0, "payload_bytes": 20}) active = queue.take() queue.put({"sequence": 1, "payload_bytes": 1}) assert queue.dropped == [{"sequence": 1, "reason": "byte-budget"}] assert queue.peak_bytes == 20 queue.release(active) queue.finish() assert queue.take() is None def test_numpy_member_is_read_incrementally_and_bounded(pilot): data = io.BytesIO() np.savez_compressed(data, points=np.arange(300, dtype="