''' Deterministic history-backfill regressions. ''' from contextlib import asynccontextmanager from functools import partial from pathlib import Path from types import SimpleNamespace from uuid import uuid4 import numpy as np from pendulum import ( datetime, from_timestamp, ) import polars as pl import pytest import tractor import trio from piker.brokers import DataUnavailable from piker.brokers.ib.api import ( _bar_load_dtype, _ohlc_dtype, ) from piker.data._sharedmem import maybe_open_shm_array from piker.data._source import def_iohlcv_fields from piker.storage.nativedb import ( NativeStorageClient, ohlc_key_map, ) from piker.tsp._history import ( _synthetic_gap_times, maybe_fill_null_segments, notify_backfill, publish_latest_frame, start_backfill, tsdb_backfill, ) def run_exhausted_backfill(get_hist) -> None: ''' Run a provider-exhaustion case without actor services. ''' async def main() -> None: with trio.fail_after(0.5): async with trio.open_nursery() as nursery: await nursery.start(partial( start_backfill, get_hist=get_hist, def_frame_duration=None, mod=SimpleNamespace(name='fake'), mkt=SimpleNamespace(fqme='x.fake'), shm=object(), timeframe=60, backfill_from_shm_index=100, backfill_from_dt=datetime(2026, 1, 2), sampler_stream=object(), backfill_until_dt=datetime(2026, 1, 1), storage=None, write_tsdb=False, )) trio.run(main) def test_data_unavailable_completes() -> None: ''' Provider exhaustion exits without orphaned completion waits. ''' async def get_hist(*args, **kwargs): raise DataUnavailable('history exhausted') run_exhausted_backfill(get_hist) def test_empty_frame_completes() -> None: ''' An empty provider frame is checked before endpoint indexing. ''' frame = np.empty( 0, dtype=np.dtype(def_iohlcv_fields), ) async def get_hist(*args, **kwargs): end_dt = kwargs['end_dt'] return frame, end_dt, end_dt run_exhausted_backfill(get_hist) def test_storage_receives_full_provider_delta() -> None: ''' SHM drops the published endpoint without truncating storage. Reverse history queries include their requested end bar, which is already the oldest sample in SHM. Publishing the overlap consumed an extra physical row and made reserved null counts disagree with elapsed timestamps. Return a frame with the overlap and prove NativeDB receives the complete provider delta while SHM receives only bars older than its published boundary. ''' frame = np.zeros( 3, dtype=np.dtype(def_iohlcv_fields), ) frame['index'] = [0, 1, 2] frame['time'] = [60, 120, 180] events: list[str] = [] class Shm: def __init__(self) -> None: self.pushed: list[np.ndarray] = [] def push(self, array, **kwargs) -> None: self.pushed.append(array.copy()) events.append('shm') class Sampler: async def send(self, msg) -> None: events.append('sampler') class Storage: def __init__(self) -> None: self.frames: list[np.ndarray] = [] async def update_ohlcv( self, fqme: str, ohlcv: np.ndarray, timeframe: int, ) -> None: self.frames.append(ohlcv.copy()) events.append('storage') shm = Shm() storage = Storage() mkt = SimpleNamespace( fqme='x.test', dst=SimpleNamespace(atype='crypto'), src=SimpleNamespace(atype='crypto_currency'), get_fqme=lambda **kwargs: 'x.test', ) async def get_hist(*args, **kwargs): return ( frame, from_timestamp(60), from_timestamp(180), ) async def main() -> None: with trio.fail_after(0.5): await start_backfill( get_hist=get_hist, def_frame_duration=None, mod=SimpleNamespace(name='fake'), mkt=mkt, shm=shm, timeframe=60, backfill_from_shm_index=2, backfill_from_dt=from_timestamp(180), sampler_stream=Sampler(), backfill_until_dt=from_timestamp(60), storage=storage, write_tsdb=True, ) trio.run(main) assert shm.pushed[0]['time'].tolist() == [60, 120] assert storage.frames[0]['time'].tolist() == [60, 120, 180] assert events == ['storage', 'shm', 'sampler'] def test_latest_frame_is_persisted_before_shm() -> None: ''' Startup persists the most-recent provider frame before publication. ''' frame = np.zeros( 2, dtype=np.dtype(def_iohlcv_fields), ) frame['time'] = [60, 120] events: list[str] = [] class Storage: async def update_ohlcv(self, *args) -> None: events.append('storage') class Shm: def push(self, array, **kwargs) -> None: events.append('shm') async def main() -> None: with trio.fail_after(0.5): await publish_latest_frame( storage=Storage(), mkt=SimpleNamespace( fqme='x.test', dst=SimpleNamespace(atype='crypto'), src=SimpleNamespace(atype='crypto_currency'), get_fqme=lambda **kwargs: 'x.test', ), shm=Shm(), array=frame, timeframe=60, ) trio.run(main) assert events == ['storage', 'shm'] def test_tsdb_prepend_uses_actual_provider_boundary( monkeypatch: pytest.MonkeyPatch, ) -> None: ''' Sparse venue history must not reserve wall-clock rows in SHM. NVDA's 60s startup converted elapsed calendar minutes between NativeDB and the latest IB frame into physical SHM offsets. IB returned only trading bars, leaving 42,720 zero rows even though both valid margins ended at the same timestamp. Model completed reverse retrieval with a provider boundary overlapping NativeDB. Exercise real ``ShmArray.push()`` indexing and prove the stored, reverse-provider, and latest frames form one contiguous lossless sequence without an explicit wall-clock offset. ''' stored = np.zeros( 3, dtype=np.dtype(def_iohlcv_fields), ) stored['time'] = [60, 120, 180] older_reverse = np.zeros( 3, dtype=np.dtype(def_iohlcv_fields), ) older_reverse['time'] = [180, 240, 300] newer_reverse = np.zeros( 3, dtype=np.dtype(def_iohlcv_fields), ) newer_reverse['time'] = [300, 360, 420] latest = np.zeros( 2, dtype=np.dtype(def_iohlcv_fields), ) latest['time'] = [420, 480] storage_frames: list[np.ndarray] = [] notifications: list[dict] = [] startup_done = trio.Event() class Storage: ''' Capture the complete provider delta written to NativeDB. ''' async def update_ohlcv( self, fqme: str, ohlcv: np.ndarray, timeframe: int, ) -> None: ''' Record one reverse provider response. ''' storage_frames.append(ohlcv.copy()) async def load( self, fqme: str, timeframe: int, ) -> tuple: ''' Return overlapping NativeDB history and its boundaries. ''' return ( stored, from_timestamp(60), from_timestamp(180), ) class Sampler: ''' Accept deterministic backfill notifications. ''' async def send(self, msg) -> None: ''' Capture notification without actor IPC. ''' notifications.append(msg) mkt = SimpleNamespace( fqme='nvda.nasdaq.ib', dst=SimpleNamespace(atype='stock'), src=SimpleNamespace(atype='fiat'), get_fqme=lambda **kwargs: 'nvda.nasdaq.ib', ) async def get_hist(*args, **kwargs): ''' Return inclusive reverse frames overlapping each SHM boundary. ''' end_dt = kwargs['end_dt'] if end_dt is None: return ( latest, from_timestamp(420), from_timestamp(480), ) end_t: float = end_dt.timestamp() reverse: np.ndarray = ( newer_reverse if end_t == 420 else older_reverse ) return ( reverse, from_timestamp(reverse['time'][0]), from_timestamp(reverse['time'][-1]), ) @asynccontextmanager async def open_history_client(mkt): ''' Yield the deterministic provider and default frame config. ''' yield get_hist, {} async def finish_null_repair(**kwargs) -> None: ''' Signal that provider and NativeDB prepends both completed. ''' startup_done.set() monkeypatch.setattr( 'piker.tsp._history.maybe_fill_null_segments', finish_null_repair, ) async def main() -> None: shm_key = f'test_sparse_prepend_{uuid4().hex}' async with tractor.open_root_actor( name=shm_key, tpt_bind_addrs=[('127.0.0.1', 0)], ): shm, opened = maybe_open_shm_array( key=shm_key, size=16, dtype=np.dtype(def_iohlcv_fields), append_start_index=12, ) assert opened storage = Storage() async with trio.open_nursery() as nursery: nursery.start_soon(partial( tsdb_backfill, mod=SimpleNamespace( name='fake', open_history_client=open_history_client, ), storemod=SimpleNamespace( ohlc_key_map=ohlc_key_map, ), storage=storage, mkt=mkt, shm=shm, timeframe=60, sampler_stream=Sampler(), )) with trio.fail_after(0.5): await startup_done.wait() nursery.cancel_scope.cancel() assert shm.array['time'].tolist() == [ 60, 120, 180, 240, 300, 360, 420, 480, ] assert not np.any(shm.array['time'] <= 0) trio.run(main) assert len(storage_frames) == 3 assert storage_frames[0]['time'].tolist() == [420, 480] assert storage_frames[1]['time'].tolist() == [300, 360, 420] assert storage_frames[2]['time'].tolist() == [180, 240, 300] assert len(notifications) == 3 def test_ib_latest_frame_round_trips_through_nativedb( tmp_path: Path, ) -> None: ''' Persist and reload IB's provider schema through history startup. ``publish_latest_frame()`` previously passed IB's index-less bars directly to NativeDB, whose durable-schema validator rejected the missing derived ``index``. IB also appends ``count`` after OHLCV, which exposed positional Polars-to-NumPy conversion to field corruption. Build the actual IB load dtype with distinct values, run the real history publication, NativeDB Parquet, and ``ShmArray`` paths, then hydrate a fresh IB buffer from storage. Assertions prove first-start publication preserves provider fields while restart maps canonical fields and leaves provider-only ``count`` at its default. ''' frame = np.zeros( 2, dtype=np.dtype(_bar_load_dtype), ) frame['time'] = [60, 120] frame['open'] = [1.1, 2.1] frame['high'] = [1.2, 2.2] frame['low'] = [1.0, 2.0] frame['close'] = [1.15, 2.15] frame['volume'] = [10, 20] frame['count'] = [3, 4] storage = NativeStorageClient(tmp_path) shm_key = f'test_ib_history_{uuid4().hex}' mkt = SimpleNamespace( fqme='mnq.cme.20260918.ib', dst=SimpleNamespace(atype='continuous_future'), src=SimpleNamespace(atype='fiat'), get_fqme=lambda **kwargs: 'mnq.cme.20260918.ib', ) async def main() -> None: with trio.fail_after(2): async with tractor.open_root_actor( name=shm_key, tpt_bind_addrs=[('127.0.0.1', 0)], ): shm, opened = maybe_open_shm_array( key=f'{shm_key}_first', size=16, dtype=np.dtype(_ohlc_dtype), append_start_index=8, ) restart_shm, restart_opened = maybe_open_shm_array( key=f'{shm_key}_restart', size=16, dtype=np.dtype(_ohlc_dtype), append_start_index=8, ) assert opened assert restart_opened await publish_latest_frame( storage=storage, mkt=mkt, shm=shm, array=frame, timeframe=60, ) loaded = await storage.read_ohlcv( mkt.fqme, timeframe=60, ) restart_shm.push( loaded, prepend=True, field_map=ohlc_key_map, ) stored = pl.read_parquet( storage.mk_path(mkt.fqme, 60) ) canonical = [ name for name, _ in def_iohlcv_fields ] canonical_schema = { name: ( pl.Int64 if field_type is int else pl.Float64 ) for name, field_type in def_iohlcv_fields } assert stored.columns == canonical assert dict(stored.schema) == canonical_schema assert list(loaded.dtype.fields) == canonical assert loaded['index'].tolist() == [0, 1] for field in canonical[1:]: assert ( loaded[field].tolist() == frame[field].tolist() ) assert ( restart_shm.array[field].tolist() == frame[field].tolist() ) assert shm.array['count'].tolist() == [3, 4] assert restart_shm.array['count'].tolist() == [0, 0] trio.run(main) def test_backfill_notification_timeout_is_bounded( monkeypatch: pytest.MonkeyPatch, ) -> None: ''' A wedged sampler can not indefinitely shield actor teardown. ''' monkeypatch.setattr( 'piker.tsp._history._notify_timeout_s', 0.01, ) class Sampler: async def send(self, msg) -> None: await trio.sleep_forever() async def main() -> None: with trio.fail_after(0.1): await notify_backfill( Sampler(), SimpleNamespace(fqme='x.test'), 60, ) trio.run(main) @pytest.mark.parametrize('epoch', [0, 1790193489]) def test_synthetic_gap_times_require_valid_right_boundary( epoch: int, ) -> None: ''' Synthetic rows must remain strictly between valid boundary bars. The live QQQ reservation held one more physical zero row than its timestamp interval could represent. Blind fill would assign the null the right boundary's timestamp, creating a duplicate. Prove aligned bounds produce cadence while an overfull segment is rejected instead of inventing non-monotonic timestamps. Repeat at a modern epoch: relative timestamp tolerance must not allow extra slots merely because epoch seconds are large. ''' aligned = _synthetic_gap_times( start_t=epoch + 60, right_t=epoch + 300, count=3, timeframe=60, ) overfull = _synthetic_gap_times( start_t=epoch + 60, right_t=epoch + 240, count=3, timeframe=60, ) trailing = _synthetic_gap_times( start_t=epoch + 60, right_t=None, count=3, timeframe=60, ) assert aligned is not None assert aligned.tolist() == [epoch + t for t in (120, 180, 240)] assert overfull is None assert trailing is None @pytest.mark.parametrize('reply', ['empty', 'nodata', 'unavailable']) def test_null_repair_falls_back_without_debug_pause( monkeypatch: pytest.MonkeyPatch, reply: str, ) -> None: ''' Empty provider repair must resolve interior zero rows locally. QQQ startup exposed 653 zero-time rows between stored history and an IB frame. The null repair queried IB while reverse backfill occupied its sole history request slot, then could enter an interactive debug pause before reaching forward-fill. Arrange an interior null run, return no provider history, and wedge sampler notification. A short timeout and fail-fast pause prove local fallback needs no interactive or IPC progress. The assertions verify only the zero rows inherit the preceding close and synthesized timestamps; both valid boundary rows remain unchanged. Exercise empty arrays and both typed no-history outcomes; each must reach the same fallback without requiring a debugger. ''' monkeypatch.setattr( 'piker.tsp._history._notify_timeout_s', 0.01, ) backing = np.zeros( 105, dtype=np.dtype(def_iohlcv_fields), ) backing['index'] = np.arange(105) frame = backing[100:105] frame['index'] = np.arange(100, 105) frame['time'] = [60, 0, 0, 0, 300] frame['open'] = [9, 0, 0, 0, 19] frame['high'] = [11, 0, 0, 0, 21] frame['low'] = [8, 0, 0, 0, 18] frame['close'] = [10, 0, 0, 0, 20] frame['volume'] = [5, 0, 0, 0, 6] class Shm: ''' Expose the test frame through the SHM repair interface. ''' def __init__(self) -> None: self._array = backing @property def array(self) -> np.ndarray: ''' Return the readable SHM view. ''' return self._array[100:105] class Sampler: ''' Model a backpressured UI notification stream. ''' async def send(self, msg) -> None: ''' Block until the bounded notifier cancels this send. ''' await trio.sleep_forever() async def get_hist(*args, **kwargs): ''' Return no provider bars for the requested null segment. ''' from piker.brokers import NoData if reply == 'nodata': raise NoData(info=kwargs) if reply == 'unavailable': raise DataUnavailable('no history') end_dt = kwargs['end_dt'] empty = np.empty( 0, dtype=np.dtype(def_iohlcv_fields), ) return empty, end_dt, end_dt async def fail_pause() -> None: ''' Reject interactive debugging in normal repair fallbacks. ''' raise AssertionError('null repair entered tractor.pause()') monkeypatch.setattr(tractor, 'pause', fail_pause) async def main() -> None: import greenback async def ensure_portal() -> None: ''' Avoid installing a greenback portal in this unit test. ''' monkeypatch.setattr( greenback, 'ensure_portal', ensure_portal, ) with trio.fail_after(0.1): await maybe_fill_null_segments( shm=Shm(), timeframe=60, get_hist=get_hist, sampler_stream=Sampler(), mkt=SimpleNamespace(fqme='qqq.nasdaq.ib'), backfill_until_dt=from_timestamp(60), ) trio.run(main) assert frame['time'].tolist() == [60, 120, 180, 240, 300] for field in ('open', 'high', 'low', 'close'): assert frame[field].tolist() == [ frame[field][0], 10, 10, 10, frame[field][-1], ] assert frame['volume'].tolist() == [5, 0, 0, 0, 6] def test_null_repair_leaves_packed_time_gap_unchanged( monkeypatch: pytest.MonkeyPatch, ) -> None: ''' A missing time interval is not an allocated zero-time segment. Cancellation during reverse backfill can leave two valid stored frames separated in time. Reload packs their rows into adjacent SHM slots. The null repair must not overwrite these valid bars with synthetic values or pretend it recovered missing trades. Use the observed MNQ gap with adjacent physical indexes; fail on any provider request and compare every field after the call. This documents the separate timestamp-gap repair still needed. ''' import greenback frame = np.zeros(2, dtype=np.dtype(def_iohlcv_fields)) frame['index'] = [100, 101] frame['time'] = [1790193489, 1790248914] for field in ('open', 'high', 'low', 'close'): frame[field] = [100, 110] frame['volume'] = [3, 7] before = frame.copy() async def ensure_portal() -> None: ''' Avoid installing a debugger portal for this unit test. ''' async def get_hist(*args, **kwargs): ''' Reject treating two valid rows as a zero-time reservation. ''' raise AssertionError('packed gap reached null repair query') monkeypatch.setattr(greenback, 'ensure_portal', ensure_portal) trio.run(partial( maybe_fill_null_segments, shm=SimpleNamespace(array=frame), timeframe=1, get_hist=get_hist, sampler_stream=object(), mkt=SimpleNamespace(fqme='mnq.cme.20261218.ib'), backfill_until_dt=from_timestamp(frame['time'][0]), )) np.testing.assert_array_equal(frame, before) @pytest.mark.parametrize('empty_reply', ['hmds_error', 'empty_bars']) def test_ib_empty_interval_preserves_older_bars_in_parquet( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, empty_reply: str, ) -> None: ''' An empty IB interval must not skip a day of available history. Run the real IB adapter, startup backfiller, SHM publication, and NativeDB writes against a simulated IB bars endpoint. A known empty interval separates recent bars from older trades. The old HMDS handler jumped a day back, so those older trades never reached SHM or parquet. Keep the legitimate empty interval and assert every supplied trade survives a fresh storage load. This covers data ingestion and persistence, not Qt rendering. ''' from ib_async import RequestError from pendulum import duration from piker.brokers.ib import feed as ib_feed epoch: int = 1790065963 requests: list[float|None] = [] finished = trio.Event() mkt = SimpleNamespace( fqme='mnq.cme.20261218.ib', dst=SimpleNamespace(atype='future'), src=SimpleNamespace(atype='fiat'), get_fqme=lambda **kwargs: 'mnq.cme.20261218.ib', get_bs_fqme=lambda **kwargs: 'mnq.cme.20261218', ) def bars(offsets) -> np.ndarray: frame = np.zeros(len(offsets), dtype=_bar_load_dtype) frame['time'] = np.asarray(offsets) + epoch for field in ('open', 'high', 'low', 'close'): frame[field] = 30000 + np.asarray(offsets) / 4 frame['volume'] = 1 frame['count'] = 1 return frame class Proxy: async def maybe_get_head_time(self, **kwargs): return from_timestamp(epoch - 864000) async def bars(self, **kwargs): end = kwargs['end_dt'] end_t = None if end is None else end.timestamp() requests.append(end_t) if end_t is None: frame = bars([4001, 4002]) elif end_t == epoch + 4001: if empty_reply == 'hmds_error': raise RequestError( 1, 162, 'HMDS query returned no data', ) frame = bars([]) elif end_t == epoch + 2001: frame = bars(list(range(1, 2002))) else: # Model the old adapter's day-skipped response. # It falls before the stored endpoint and is dropped. assert end_t == epoch + 4001 - 86400 frame = bars([-82400, -82399]) native = [ SimpleNamespace(date=from_timestamp(timestamp)) for timestamp in frame['time'] ] return native, frame, duration(seconds=2000) @asynccontextmanager async def fake_data_client(): yield Proxy() class Sampler: async def send(self, msg) -> None: pass # Preserve the real null-repair pass; signal when startup ends. from piker.tsp import _history repair = _history.maybe_fill_null_segments async def finish_repair(**kwargs) -> None: await repair(**kwargs) finished.set() monkeypatch.setattr( ib_feed, 'open_data_client', fake_data_client, ) monkeypatch.setattr( _history, 'maybe_fill_null_segments', finish_repair, ) expected = bars(list(range(2002)) + [4001, 4002]) async def main() -> None: storage = NativeStorageClient(tmp_path) await storage.update_ohlcv(mkt.fqme, bars([0, 1]), 1) key = f'test_ib_request_bounds_{uuid4().hex}' with trio.fail_after(10): async with tractor.open_root_actor( name=key, tpt_bind_addrs=[('127.0.0.1', 0)], ): shm, opened = maybe_open_shm_array( key=key, size=5000, dtype=np.dtype(_ohlc_dtype), append_start_index=4500, ) assert opened async with trio.open_nursery() as nursery: nursery.start_soon(partial( tsdb_backfill, mod=SimpleNamespace( name='ib', open_history_client=( ib_feed.open_history_client ), ), storemod=SimpleNamespace( ohlc_key_map=ohlc_key_map, ), storage=storage, mkt=mkt, shm=shm, timeframe=1, sampler_stream=Sampler(), )) await finished.wait() nursery.cancel_scope.cancel() # Read actual parquet through a new client to avoid # accidentally validating only NativeDB's cache. fresh = NativeStorageClient(tmp_path) loaded, _, _ = await fresh.load(mkt.fqme, 1) for field in ( 'time', 'open', 'high', 'low', 'close', 'volume', ): np.testing.assert_array_equal( loaded[field], expected[field], err_msg=f'parquet field: {field}', ) np.testing.assert_array_equal( shm.array[field], expected[field], err_msg=f'SHM field: {field}', ) assert requests == [ None, epoch + 4001, epoch + 2001, ] # Exactly the simulated closure remains. Do not # manufacture bars to force timestamp continuity. gaps = np.diff(loaded['time']) assert gaps[gaps > 1].tolist() == [2000] trio.run(main) @pytest.mark.parametrize('case', [ 'aligned', 'overfull', 'unbounded', 'occupied', ]) def test_copy_close_to_null_rows_preserves_other_fields( case: str, ) -> None: ''' Synthetic repair must be isolated from SHM and provider I/O. Reserve two rows between valid bars at a modern epoch. Only a fully null, correctly bounded run may receive the preceding close and generated timestamps. Reject an occupied row, a duplicate right timestamp, or a missing right margin without any partial mutation. Distinct volume/count values prove this helper copies neither the seed's activity nor entire records. ''' from piker.tsp._history import _copy_close_to_null_rows frame: np.ndarray = np.zeros(4, dtype=np.dtype(_ohlc_dtype)) epoch: int = 1790193489 frame['index'] = np.arange(100, 104) frame['time'] = [epoch, 0, 0, epoch + 3] frame['close'] = [123, 0, 0, 456] frame['volume'] = [8, 2, 3, 9] frame['count'] = [4, 5, 6, 7] right_t: float|None = float(frame['time'][-1]) if case == 'overfull': right_t -= 1 elif case == 'unbounded': right_t = None elif case == 'occupied': frame['time'][1] = epoch + 1 before: np.ndarray = frame.copy() copied: int = _copy_close_to_null_rows( frame[1:3], frame[0], right_t, 1, ) expected: np.ndarray = before.copy() if case == 'aligned': expected['time'][1:3] = [epoch + 1, epoch + 2] for field in ('open', 'high', 'low', 'close'): expected[field][1:3] = 123 assert copied == 2 else: assert copied == 0 np.testing.assert_array_equal(frame, expected)