diff --git a/ai/prompt-io/codex/20260924T165757Z_1da9b760_prompt_io.md b/ai/prompt-io/codex/20260924T165757Z_1da9b760_prompt_io.md new file mode 100644 index 00000000..92ee6915 --- /dev/null +++ b/ai/prompt-io/codex/20260924T165757Z_1da9b760_prompt_io.md @@ -0,0 +1,51 @@ +--- +model: gpt-6 (exact variant not exposed) +service: codex +timestamp: 2026-09-24T16:57:57.342784+00:00 +git_ref: 1da9b760 +scope: code +substantive: true +raw_file: 20260924T165757Z_1da9b760_prompt_io.raw.md +--- + +## Prompt + +Document and test the existing null-segment repair, including its +flat-bar fallback, while distinguishing packed timestamp gaps. + +## Response summary + +Documented SHM-only repair, provider placement assumptions, fallback +fields, notifications, and concurrency limits. Added boundary and +no-history coverage and fixed epoch-dependent comparison tolerance +which could synthesize a duplicate timestamp at the right margin. + +## Files changed + +- piker/tsp/_history.py: helper contract and absolute tolerance. +- tests/test_history_backfill.py: four additional parameter/case runs. + +## Human edits + +The human identified cancellation as a source of stored gaps, +correctly described packed indexes and insertion costs, pointed to +maybe_fill_null_segments(), and requested immediate documentation +and tests. No direct source edits observed during this change; +the patch remains uncommitted for review. + +## Close-copy factoring follow-up + +The human requested stronger internal types, a reusable close-copy +helper, searchable comments for copying close and synthesizing +stamps, and research into historical venue closure queries. +Extracted _copy_close_to_null_rows(), preserving non-OHLC fields and +rejecting occupied slots before mutation. Added explicit callback, +array, segment-index and date annotations. Preserved the existing +provenance TODO and suggested avoiding volume sentinels. +Four helper cases cover successful fill and unchanged rejected +inputs. Backfill suite: 19 passed in 2.01s; Ruff and diff check pass. + +IB documents historical SCHEDULE requests separately from contract +hours. The installed ib_async exposes reqHistoricalScheduleAsync. +No live schedule request or venue behavior change was performed. +Source: https://interactivebrokers.github.io/tws-api/historical_bars.html diff --git a/ai/prompt-io/codex/20260924T165757Z_1da9b760_prompt_io.raw.md b/ai/prompt-io/codex/20260924T165757Z_1da9b760_prompt_io.raw.md new file mode 100644 index 00000000..c84c7063 --- /dev/null +++ b/ai/prompt-io/codex/20260924T165757Z_1da9b760_prompt_io.raw.md @@ -0,0 +1,48 @@ +--- +model: gpt-6 (exact variant not exposed) +service: codex +timestamp: 2026-09-24T16:57:57.342784+00:00 +git_ref: 1da9b760 +diff_cmd: git diff 1da9b760 -- piker/tsp/_history.py tests/test_history_backfill.py +--- + +## Prompt + +The human requested documentation and tests for +maybe_fill_null_segments(), identifying its flat OHLC fallback +and asking how packed SHM could accept later recovered bars. + +## Code output + +> `git diff 1da9b760 -- piker/tsp/_history.py` + +> `git diff 1da9b760 -- tests/test_history_backfill.py` + +Document the existing slot-based repair and its limitations. +Extend fallback coverage to NoData and DataUnavailable; verify +the observed packed MNQ gap remains unchanged. Reproduce the +relative-tolerance failure at a modern epoch, then disable relative +tolerance in the synthetic timestamp boundary check. + +## Validation + +Before the fix: modern-epoch boundary case failed, zero-epoch +case passed. The overfull synthetic run duplicated the right bar. +After the fix: 15 backfill tests passed in 1.99s. + +## Close-copy factoring follow-up + +The human requested stronger internal types, a reusable close-copy +helper, searchable comments for copying close and synthesizing +stamps, and research into historical venue closure queries. +Extracted _copy_close_to_null_rows(), preserving non-OHLC fields and +rejecting occupied slots before mutation. Added explicit callback, +array, segment-index and date annotations. Preserved the existing +provenance TODO and suggested avoiding volume sentinels. +Four helper cases cover successful fill and unchanged rejected +inputs. Backfill suite: 19 passed in 2.01s; Ruff and diff check pass. + +IB documents historical SCHEDULE requests separately from contract +hours. The installed ib_async exposes reqHistoricalScheduleAsync. +No live schedule request or venue behavior change was performed. +Source: https://interactivebrokers.github.io/tws-api/historical_bars.html diff --git a/piker/tsp/_history.py b/piker/tsp/_history.py index d0d6bd5f..87b0c613 100644 --- a/piker/tsp/_history.py +++ b/piker/tsp/_history.py @@ -38,6 +38,7 @@ from pprint import pformat from string import hexdigits from types import ModuleType from typing import ( + Awaitable, Callable, Generator, Literal, @@ -87,7 +88,6 @@ from piker.storage import TimeseriesNotFound from ._anal import ( get_null_segs, iter_null_segs, - Frame, # codec-ish np2pl as np2pl, @@ -225,12 +225,63 @@ def _synthetic_gap_times( if not np.isclose( right_t, start_t + (count + 1)*timeframe, + # Epoch magnitude must not widen allowed boundary error. + rtol=0, + atol=1e-6, ): return None return times +def _copy_close_to_null_rows( + gap: np.ndarray, + seed: np.void, + right_t: float|None, + timeframe: float, + +) -> int: + ''' + Copy preceding CLOSE into flat OHLC and SYNTHESIZE TIMESTAMPS. + + Mutate only an entirely zero-time `gap` view, bounded by `seed` + and `right_t`. Return the number of synthetic rows, or zero + without mutation when the bounds cannot accommodate the slots. + Preserve indexes and non-OHLC fields, including volume/count. + This synchronous operation performs no I/O or SHM publication. + + These are synthetic placeholders, not provider observations. + Callers own provenance tracking and subscriber notification. + + ''' + start_t: float = float(seed['time']) + if ( + not len(gap) + or + start_t <= 0 + or + not np.all(gap['time'] == 0) + ): + return 0 + + gap_times: np.ndarray|None = _synthetic_gap_times( + start_t, + right_t, + len(gap), + timeframe, + ) + if gap_times is None: + return 0 + + # SYNTHETIC FLAT BARS: COPY PRECEDING CLOSE into O/H/L/C. + # SYNTHESIZE TIMESTAMPS only for the existing zero-time slots. + # No volume/count fabrication and no physical row insertion. + close: float = float(seed['close']) + gap[['open', 'high', 'low', 'close']] = close + gap['time'] = gap_times + return len(gap) + + async def shm_push_in_between( shm: ShmArray, to_push: np.ndarray, @@ -276,7 +327,9 @@ async def shm_push_in_between( async def maybe_fill_null_segments( shm: ShmArray, timeframe: float, - get_hist: Callable, + get_hist: Callable[..., Awaitable[ + tuple[np.ndarray, datetime, datetime] + ]], sampler_stream: tractor.MsgStream, mkt: MktPair, backfill_until_dt: datetime, @@ -284,19 +337,67 @@ async def maybe_fill_null_segments( task_status: TaskStatus[None] = trio.TASK_STATUS_IGNORED, ) -> None: + ''' + Repair existing zero-time SHM rows, then try flat-bar fallback. + + This operates on physical slots already inside `shm.array`. + `get_null_segs()` finds rows whose time is zero; a timestamp + jump between two nonzero rows does not trigger a query. No rows + are inserted, shifted, or reindexed, and no storage is updated. + + Signal task startup, snapshot the readable frame, and request + provider history using each null segment's surrounding dates. + Provider replies are written backward from the segment's right + boundary through `shm_push_in_between()`, without moving the + readable first index. This legacy placement assumes the reply + fits the existing slot layout; it is not a general timestamp + merge and does not limit writes to only the null rows. + + An empty reply, NoData, DataUnavailable, or a reply beginning + before `backfill_until_dt` stops further provider queries for + this pass. Cancellation and other provider errors propagate. + + Rescan the live frame for remaining zero-time runs. For each + bounded run accepted by `_synthetic_gap_times()`, copy the + preceding valid bar's close into every OHLC field and generate + timestamps at `timeframe` spacing. Preserve both boundary bars, + physical indexes, and all other fields (including volume). + Leading/trailing runs and incompatible slot counts stay null. + + Flat bars are synthetic display fallback, not recovered trades + or proof that the venue was closed. They currently have no + explicit provenance flag. Notify subscribers after provider + writes and after synthetic repair via `notify_backfill()`. + Callers must coordinate SHM writers and provider access; this + helper does not lock either across its await points. + + ''' task_status.started() - frame: Frame = shm.array.copy() + frame: np.ndarray = shm.array.copy() # TODO, put in parent task/daemon root! import greenback await greenback.ensure_portal() - null_segs: tuple|None = get_null_segs( + null_segs: tuple[ + list[list[int]], np.ndarray, np.ndarray + ]|None = get_null_segs( frame, period=timeframe, ) + absi_start: int + absi_end: int + fi_start: int|None + fi_end: int + start_t: float|None + end_t: float + start_dt: datetime|None + end_dt: datetime + array: np.ndarray + next_start_dt: datetime + next_end_dt: datetime for ( absi_start, absi_end, fi_start, fi_end, @@ -387,8 +488,8 @@ async def maybe_fill_null_segments( ) # RECHECK for more null-gaps - frame: Frame = shm.array - null_segs: tuple | None = get_null_segs( + frame = shm.array + null_segs = get_null_segs( frame, period=timeframe, ) @@ -397,6 +498,9 @@ async def maybe_fill_null_segments( and len(null_segs[-1]) ): + iabs_slices: list[list[int]] + iabs_zero_rows: np.ndarray + _zero_t: np.ndarray ( iabs_slices, iabs_zero_rows, @@ -412,12 +516,6 @@ async def maybe_fill_null_segments( # stretching the y-axis.. # array: np.ndarray = shm.array # zeros = array[array['low'] == 0] - ohlc_fields: list[str] = [ - 'open', - 'high', - 'low', - 'close', - ] zero_indexes: np.ndarray = np.asarray(iabs_zero_rows) split_at: np.ndarray = ( @@ -431,6 +529,7 @@ async def maybe_fill_null_segments( ) repaired: int = 0 frame_start: int = int(frame['index'][0]) + zero_group: np.ndarray for zero_group in zero_groups: istart: int = int(zero_group[0]) istop: int = int(zero_group[-1]) + 1 @@ -453,7 +552,6 @@ async def maybe_fill_null_segments( # Fill only the null rows, preserving both valid margins. gap: np.ndarray = shm._array[istart:istop] - cls: float = seed['close'] frame_end: int = int(frame['index'][-1]) right_t: float|None = ( @@ -461,29 +559,28 @@ async def maybe_fill_null_segments( if istop <= frame_end else None ) - gap_times: np.ndarray|None = _synthetic_gap_times( - start_t, + # TODO: how can we mark this range as being a gap tho? + # -[ ] maybe pg finally supports nulls in ndarray to + # show empty space somehow? + # -[ ] we could put a special value in the vlm or + # another col/field to denote? + # Prefer explicit synthetic-row provenance over a volume + # sentinel: zero volume alone does not identify a repair. + # COPY CLOSE / SYNTHESIZE TIMESTAMPS fallback lives here. + copied: int = _copy_close_to_null_rows( + gap, + seed, right_t, - len(gap), timeframe, ) - if gap_times is None: + if not copied: log.warning( f'Can not safely forward-fill SHM nulls:\n' f'left={start_t} right={right_t} ' f'rows={len(gap)} period={timeframe}\n' ) continue - - # TODO: how can we mark this range as being a gap tho? - # -[ ] maybe pg finally supports nulls in ndarray to - # show empty space somehow? - # -[ ] we could put a special value in the vlm or - # another col/field to denote? - gap[ohlc_fields] = cls - - gap['time'] = gap_times - repaired += len(gap) + repaired += copied # TODO: reimpl using the new `.ui._remote_ctl` ctx # ideally using some kinda decent diff --git a/tests/test_history_backfill.py b/tests/test_history_backfill.py index d0b19ead..b756ad64 100644 --- a/tests/test_history_backfill.py +++ b/tests/test_history_backfill.py @@ -548,7 +548,10 @@ def test_backfill_notification_timeout_is_bounded( trio.run(main) -def test_synthetic_gap_times_require_valid_right_boundary() -> None: +@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. @@ -557,35 +560,39 @@ def test_synthetic_gap_times_require_valid_right_boundary() -> None: 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=60, - right_t=300, + start_t=epoch + 60, + right_t=epoch + 300, count=3, timeframe=60, ) overfull = _synthetic_gap_times( - start_t=60, - right_t=240, + start_t=epoch + 60, + right_t=epoch + 240, count=3, timeframe=60, ) trailing = _synthetic_gap_times( - start_t=60, + start_t=epoch + 60, right_t=None, count=3, timeframe=60, ) assert aligned is not None - assert aligned.tolist() == [120, 180, 240] + 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. @@ -594,12 +601,14 @@ def test_null_repair_falls_back_without_debug_pause( 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 an empty provider frame, and + 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( @@ -654,6 +663,12 @@ def test_null_repair_falls_back_without_debug_pause( 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, @@ -708,6 +723,57 @@ def test_null_repair_falls_back_without_debug_pause( 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, @@ -863,3 +929,52 @@ def test_ib_empty_interval_preserves_older_bars_in_parquet( 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)