diff --git a/ai/prompt-io/codex/20260924T185456Z_1da9b760_schedule_prompt_io.md b/ai/prompt-io/codex/20260924T185456Z_1da9b760_schedule_prompt_io.md new file mode 100644 index 00000000..283b8783 --- /dev/null +++ b/ai/prompt-io/codex/20260924T185456Z_1da9b760_schedule_prompt_io.md @@ -0,0 +1,91 @@ +--- +model: gpt-6 (exact variant not exposed) +service: codex +timestamp: 2026-09-24T18:54:56.893490+00:00 +git_ref: 1da9b760 +scope: code +substantive: true +raw_file: 20260924T185456Z_1da9b760_schedule_prompt_io.raw.md +--- + +## Prompt + +The human requested: first freeze a commit-plan for the existing +close-copy/null-repair changes, then implement ib_async historical +schedule support with tests, use it for datad.ib history-buffer gap +spacing, and plan a common on-demand segment interface for CLI and +chart selection. Avoid reindexing in this first pass and record +fallback insertion only as a TODO. + +## Response summary + +Implemented historical schedule normalization, bounded startup +reservations, explicit actor-local synthetic ranges, and two read +RPC functions. Added deterministic adapter tests and real SHM/parquet +startup coverage. Wrote the segment-backfill interface plan. + +This entry covers the new schedule work only. The earlier frozen +null-repair plan has its own provenance. The shared _history.py file +contains changes from both scopes until the earlier commit is made. + +## Human edits + +The human identified cancellation during backfill as the reproducible +trigger, explained the packed-index insertion tradeoff, requested +reuse of venue closure knowledge, and chose placeholders to avoid +expensive rewrites. They required explicit placeholder annotation, +proposed later graphics compression, and asked to defer insertion +fallbacks. They also requested the existing null-repair commit plan +before the new schedule work. These are human-directed architecture +and scope choices; no acceptance or task completion state is inferred. + +## Validation and limits + +55 targeted tests passed; Ruff and diff checks passed. No live IB +qualification was performed. Reservations are synthetic, not recovered +observations. Repair execution, persistence of unfinished requests, +chart compression, FSP handling, export filtering, and request-level +IB cancellation remain follow-up work documented in the plan. + +## Files changed + +- `piker/brokers/ib/venues.py` +- `piker/brokers/ib/api.py` +- `piker/brokers/ib/feed.py` +- `piker/tsp/_reserve.py` +- `piker/tsp/_history.py` +- `tests/providers/test_ib_schedule.py` +- `tests/providers/test_history_reservations.py` +- `docs/history_segment_backfill.rst` + +## Pre-commit module review + +The human requested placing historical schedule normalization in the +existing venues.py module, naming the allocation module _reserve.py, +and adding the missing project license headers. Applied those moves, +updated imports and test patch targets, and added the standard AGPL +headers to the new allocation module and both new test modules. + +## Local tuicr review follow-up + +The human supplied eleven comments requesting a configurable newest +gap count, caller and metadata references, explanations of guards +and provider-specific limits, multiline boolean style, capacity +fallback rationale, an explanation of numpy split boundaries, and +clearer test cases, fixture types, and module docstrings. + +Added newest_gaps=32, with zero disabling queries and negative values +rejected. Documented tsdb_backfill ownership and the local/RPC range +accessors. Explained recoverable overflow and half-open synthetic +range grouping. Typed tmp_path as pathlib.Path and monkeypatch as +pytest.MonkeyPatch. Added four request-budget cases. + +Verification: reservation tests 11 passed; schedule tests 7 passed; +Ruff and git diff --check passed. Changes remain unstaged for review. +No forge review or CI state was queried; this was local feedback. + +## Provider-test layout review + +The human requested grouping both new history/schedule test modules +under tests/providers. Moved both files without changing test logic; +they retain the shared tests/conftest.py fixtures. diff --git a/ai/prompt-io/codex/20260924T185456Z_1da9b760_schedule_prompt_io.raw.md b/ai/prompt-io/codex/20260924T185456Z_1da9b760_schedule_prompt_io.raw.md new file mode 100644 index 00000000..41f34b42 --- /dev/null +++ b/ai/prompt-io/codex/20260924T185456Z_1da9b760_schedule_prompt_io.raw.md @@ -0,0 +1,100 @@ +--- +model: gpt-6 (exact variant not exposed) +service: codex +timestamp: 2026-09-24T18:54:56.893490+00:00 +git_ref: 1da9b760 +diff_cmd: git diff HEAD~1..HEAD +--- + +> `git diff HEAD~1..HEAD -- piker/brokers/ib/_schedule.py` + +> `git diff HEAD~1..HEAD -- piker/brokers/ib/api.py` + +> `git diff HEAD~1..HEAD -- piker/brokers/ib/feed.py` + +> `git diff HEAD~1..HEAD -- piker/tsp/_reservations.py` + +> `git diff HEAD~1..HEAD -- piker/tsp/_history.py` + +> `git diff HEAD~1..HEAD -- tests/test_ib_schedule.py` + +> `git diff HEAD~1..HEAD -- tests/test_history_reservations.py` + +> `git diff HEAD~1..HEAD -- docs/history_segment_backfill.rst` + +## Generated output + +Added Client.history_schedule() over the installed ib_async API, +normalizing dated sessions to UTC with explicit response coverage. +Reject ambiguous DST times, reversed sessions, and overlapping +sessions. Match the existing extended-hours bars policy. + +Added startup-only schedule reservations to tsdb_backfill before +prepending older stored history. Reserve open-session sample slots +only where the reply covers the full candidate gap. Return synthetic +UTC ranges separately from OHLCV; copy preceding close into flat +placeholders without writing those placeholders to parquet. +Unknown coverage and capacity failures retain packed time gaps. + +Added datad.ib schedule and reservation metadata RPC functions and +a proposed shared CLI/chart segment request protocol. The repair +queue, durable unresolved intervals, chart compression, FSP handling, +and coordinated insertion remain deferred. No live data was repaired. + +Verification is complete: 55 targeted tests pass, including the +new schedule and real SHM/parquet tests. I also added checks for +reversed session bounds and the 25-hour day at the autumn DST +change. + +The segment API plan uses UTC ranges for both CLI and chart +requests, a provider-owned queue, durable unfinished intervals, +and replacement of reserved slots without reindexing. Live IB +schedule coverage still needs qualification. + +## Limits + +Only stored-history gaps are reserved before prepend; published rows +are never shifted. Newest 32 gaps, at most 31 days each, are queried +serially with a 10-second schedule timeout. This can delay history +publication. Underlying IB request cancellation is still a TODO. +Manual whole-SHM export requires synthetic-range exclusion before +it can preserve provider-only storage. The chart currently displays +flat placeholders, and synthetic-aware FSP treatment is deferred. + +Tests: history backfill 19, IB history 17, method proxy 5, schedule 7, +reservations 7. Sandbox local-socket failures were rerun with approval. +Ruff and git diff --check passed. No live broker requests were made. + +## Module-layout review response + +Schedule normalization now lives in piker.brokers.ib.venues beside +the existing session helpers. The allocation module is tsp._reserve. +Added the project AGPL headers to the new module and test files. +Earlier diff pointers record the original generated paths; use +the following paths for the reviewed layout: + +> `git diff HEAD~1..HEAD -- piker/brokers/ib/venues.py` + +> `git diff HEAD~1..HEAD -- piker/tsp/_reserve.py` + +## Local review response + +The fixes pass all 18 tests in the two affected files. +The new limit tests cover zero, negative, one-gap, and two-gap +budgets, including newest-first selection. + +Capacity overflow remains a warning: no partial reservation is +published for that gap, and the original observations remain +available. The comments now explain that fallback explicitly. + +Diff pointers above include the review edits to tsp._reserve, both +test modules, and the interface document. This review did not query +a forge or publish replies. + +## Provider-test paths + +Moved both new test modules under tests/providers as requested. + +> `git diff HEAD~1..HEAD -- tests/providers/test_ib_schedule.py` + +> `git diff HEAD~1..HEAD -- tests/providers/test_history_reservations.py` diff --git a/docs/history_segment_backfill.rst b/docs/history_segment_backfill.rst new file mode 100644 index 00000000..92a535da --- /dev/null +++ b/docs/history_segment_backfill.rst @@ -0,0 +1,113 @@ +History segment backfill: first-pass interface plan +================================================= + +Status and scope +---------------- + +The first implementation queries dated IB schedules and reserves +synthetic samples while preparing older stored history for SHM. +It does not move published rows, implement automatic segment repair, +or change chart rendering. Synthetic ranges are actor-local metadata; +parquet remains provider observations and is rescanned on restart. + +The existing null-row repair is not the segment-repair implementation. +It assumes allocated zero-time slots. Schedule reservations instead +have valid synthetic timestamps and explicit external provenance. + +Schedule queries +---------------- + +``datad.ib`` exposes ``feed.get_history_schedule(fqme, start, end)`` +for UTC epoch bounds. The response includes actual coverage, timezone, +regular-hours policy, and half-open trading sessions. This wraps +``Client.history_schedule()`` through the existing asyncio method proxy. +The first pass limits a query to 31 days and times out after 10 seconds. +Historical queries must match bars' ``useRTH=False`` selection. + +Do not extrapolate sessions outside response coverage. Empty schedules, +ambiguous DST times, failures, and partial coverage are unknown for +allocation purposes. Missing contract/session support is not a closure. +A successful schedule establishes possible trading slots, not actual +bar counts or proof of trades. Live MNQ coverage remains to be verified. + +Startup reservations +-------------------- + +After reverse retrieval and before prepending stored history, inspect +its newest 32 timestamp gaps within the retained SHM capacity by +default. ``reserve_history_gaps(newest_gaps=32)`` makes that request +budget configurable; zero disables reservations. Query schedules +serially through the history client's proxy. Reserve only +covered open-session slots, preserving venue closures as time gaps. +Large or unsupported ranges remain unresolved. Capacity bounds include +synthetic samples; oldest samples may be clipped as with ordinary SHM +prepending. No published provider or realtime row is shifted. + +``feed.get_history_reservations(fqme, timeframe)`` returns half-open UTC +ranges of synthetic rows for this datad generation. Consumers must use +this provenance rather than zero volume or flat prices. Index positions +are not durable identity. Metadata is published with the synchronous SHM +prepend and clipped to its visible timestamps. + +The default chart still displays flat placeholders. Compression and +synthetic-aware FSP treatment are follow-up consumers of this metadata. +Manual whole-SHM export must exclude synthetic ranges before writing +provider history; automatic startup writes do not include reservations. + +On-demand request contract (proposed) +------------------------------------ + +A datad context should accept ``fqme``, ``timeframe``, ``start``, ``end``, +``request_id``, and ``reason``. Bounds are half-open UTC timestamps; +reasons include interrupted startup, detected gap, and operator selection. +CLI and chart selection submit the same request, with no client-local +array indexes. The actor validates market identity and sampling grid. + +Proposed events are accepted, progress, completed, partial, failed, and +cancelled. Progress carries requested/received bounds, rows persisted, +slots replaced, remaining synthetic ranges, and failure details. +``completed`` means every requested interval received a classified +provider outcome, not that every scheduled slot contained a trade. + +A provider-owned queue deduplicates overlapping requests and serializes +IB history calls. Persist unresolved requests and per-interval outcomes +before fetching so actor/client cancellation cannot erase pending work. +Distinguish genuine empty replies from timeouts and permission failures. +Closing a progress subscription must not implicitly discard queued work; +explicit cancellation policy belongs in the request protocol. + +Merge real rows into NativeDB first, then replace matching synthetic +SHM slots by timestamp without changing indexes. Remove provenance only +after successful publication. Revalidate bounds and generation before +writes, coordinate sampler/FSP readers, and broadcast a repair interval +for downstream cache invalidation. Never use legacy null repair to place +an arbitrary provider frame backward over already valid observations. + +CLI proposal: ``piker store backfill FQME --timeframe N --start UTC +--end UTC`` with a preview showing schedules, unknown ranges, and slots. +Chart proposal: translate selected x coordinates through the current +view's timestamp mapping, preview the interval, and subscribe to the +same progress context. Compression must retain inverse timestamp/index +mapping for cursor, selection, and annotation behavior. + +Deferred work +------------- + +TODO: replace synthetic rows using the common segment request context. +TODO: durable unresolved/checked interval tracking and bounded retries. +TODO: automatic scan worker separate from samplerd's live sampling loop. +TODO: chart compression and explicit FSP handling of placeholders. +TODO: prevent operational whole-SHM exports from persisting placeholders. +TODO: cancel the underlying IB schedule request on timeout/disconnect. +TODO: batch/cache session queries with coverage and contract identity. +TODO: coordinated insertion only if schedule-driven reservation proves +insufficient; first pass reports unknown/capacity cases without reindexing. + +Validation +---------- + +Use fake schedule responses for dated timezone conversion, DST ambiguity, +partial coverage, closures, empty/error results, and capacity limits. +An isolated real-SHM/NativeDB restart test must show that reserved rows +are marked synthetic, closures stay compressed, and parquet never gains +invented trades. Live gateway qualification is separate from these tests. diff --git a/piker/brokers/ib/api.py b/piker/brokers/ib/api.py index f9e736cd..d9edc2e0 100644 --- a/piker/brokers/ib/api.py +++ b/piker/brokers/ib/api.py @@ -350,6 +350,48 @@ class Client: apiOnly=False, ) + async def history_schedule( + self, + fqme: str, + start_dt: datetime, + end_dt: datetime, + use_rth: bool = False, + timeout_s: float = 10, + ) -> dict: + ''' + Query dated trading sessions for a bounded historical range. + + Return UTC coverage separately from sessions. Missing coverage + is unknown, never evidence that the venue was closed. + Match bars()'s extended-hours policy by default. + + ''' + from math import ceil + from .venues import normalize_schedule + + if start_dt.tzinfo is None or end_dt.tzinfo is None: + raise ValueError('Schedule bounds must be timezone-aware') + seconds: float = ( + end_dt.astimezone(UTC) - start_dt.astimezone(UTC) + ).total_seconds() + if not 0 < seconds <= 31 * 86400: + raise ValueError('Schedule range must be within 31 days') + contract: Contract = (await self.find_contracts(fqme))[0] + schedule = await asyncio.wait_for( + self.ib.reqHistoricalScheduleAsync( + contract, + numDays=max(1, ceil(seconds / 86400)), + endDateTime=end_dt.astimezone(UTC).strftime( + '%Y%m%d-%H:%M:%S' + ), + useRTH=use_rth, + ), + timeout=timeout_s, + ) + # TODO: explicit cancellation of the underlying IB request + # requires a request-id exposing ib_async schedule wrapper. + return normalize_schedule(schedule, use_rth) + async def bars( self, fqme: str, diff --git a/piker/brokers/ib/feed.py b/piker/brokers/ib/feed.py index ff5b7576..2aea8b78 100644 --- a/piker/brokers/ib/feed.py +++ b/piker/brokers/ib/feed.py @@ -317,6 +317,18 @@ async def open_history_client( last_dt, ) + async def query_schedule(left: float, right: float) -> dict: + ''' + Request dated sessions through the existing IB proxy. + + ''' + return await proxy.history_schedule( + fqme=fqme, + start_dt=from_timestamp(left), + end_dt=from_timestamp(right), + use_rth=False, + ) + # TODO: it seems like we can do async queries for ohlc # but getting the order right still isn't working and I'm not # quite sure why.. needs some tinkering and probably @@ -325,6 +337,7 @@ async def open_history_client( yield ( get_hist, { + 'query_schedule': query_schedule, 'erlangs': 1, # max conc reqs 'rate': 3, # max req rate 'frame_types': { # expected frame sizes @@ -1334,3 +1347,33 @@ async def stream_quotes( # ugh, clear ticks since we've consumed them # ticker.ticks = [] # last = time.time() + + +async def get_history_schedule( + fqme: str, + start: float, + end: float, +) -> dict: + ''' + Query historical sessions on datad.ib using UTC epoch bounds. + + ''' + async with open_data_client() as proxy: + return await proxy.history_schedule( + fqme=fqme.removesuffix('.ib'), + start_dt=from_timestamp(start), + end_dt=from_timestamp(end), + use_rth=False, + ) + + +async def get_history_reservations( + fqme: str, + timeframe: float, +) -> list[tuple[float, float]]: + ''' + Report synthetic history ranges for this datad generation. + + ''' + from piker.tsp._reserve import reservation_ranges + return reservation_ranges(fqme, timeframe) diff --git a/piker/brokers/ib/venues.py b/piker/brokers/ib/venues.py index 85969707..7c529978 100644 --- a/piker/brokers/ib/venues.py +++ b/piker/brokers/ib/venues.py @@ -30,8 +30,11 @@ from datetime import ( # noqa from typing import ( Iterator, TYPE_CHECKING, + TypedDict, ) +from zoneinfo import ZoneInfo + import exchange_calendars as xcals from exchange_calendars.errors import ( InvalidCalendarName, @@ -51,6 +54,7 @@ log = get_logger(__name__) if TYPE_CHECKING: from ib_async import ( TradingSession, + HistoricalSchedule, Contract, ContractDetails, ) @@ -64,6 +68,72 @@ if TYPE_CHECKING: ) +class HistorySchedule(TypedDict): + ''' + IPC-safe UTC coverage and half-open session intervals. + + Coverage describes the server's response, not the requested + interval. Consumers must reject uncovered or ambiguous dates. + + ''' + start: float + end: float + sessions: list[tuple[float, float]] + timezone: str + use_rth: bool + + +def normalize_schedule( + schedule: HistoricalSchedule, + use_rth: bool, +) -> HistorySchedule: + ''' + Parse dated session boundaries and reject DST ambiguity. + + Never repeat today's trading hours across historical dates. + Ambiguous/nonexistent local times need provider clarification. + + ''' + zone: ZoneInfo = ZoneInfo(schedule.timeZone) + + def stamp(value: str) -> float: + ''' + Convert one unambiguous IB wall time into epoch seconds. + + ''' + local: datetime = datetime.strptime(value, '%Y%m%d-%H:%M:%S') + first: datetime = local.replace(tzinfo=zone, fold=0) + second: datetime = local.replace(tzinfo=zone, fold=1) + if first.utcoffset() != second.utcoffset(): + raise ValueError(f'Ambiguous schedule time: {value}') + return first.timestamp() + + start: float = stamp(schedule.startDateTime) + end: float = stamp(schedule.endDateTime) + if start >= end: + raise ValueError('Invalid schedule coverage') + sessions: list[tuple[float, float]] = [] + for session in schedule.sessions: + left: float = stamp(session.startDateTime) + right: float = stamp(session.endDateTime) + if left >= right: + raise ValueError('Invalid schedule session') + left = max(start, left) + right = min(end, right) + if left < right: + sessions.append((left, right)) + sessions.sort() + if any(a[1] > b[0] for a, b in zip(sessions, sessions[1:])): + raise ValueError('Overlapping schedule sessions') + return HistorySchedule( + start=start, + end=end, + sessions=sessions, + timezone=schedule.timeZone, + use_rth=use_rth, + ) + + def is_expired( con_deats: ContractDetails, ) -> bool: diff --git a/piker/tsp/_history.py b/piker/tsp/_history.py index 87b0c613..66a9e6bc 100644 --- a/piker/tsp/_history.py +++ b/piker/tsp/_history.py @@ -1506,11 +1506,42 @@ async def tsdb_backfill( # establishes the actual earliest sample, then prepend # storage directly beside it without reserving closure # rows. + from ._reserve import ( + _reservations, + reserve_history_gaps, + ) + synthetic: list[tuple[float, float]] = [] + query_schedule = config.get('query_schedule') + if ( + query_schedule is not None + and + 'time' in tsdb_history.dtype.names + ): + tsdb_history, synthetic = await reserve_history_gaps( + tsdb_history, + timeframe, + query_schedule, + max_rows=max(0, int(shm._first.value)), + ) + # Startup-only placement: the provider prefix is + # already published, but these older rows are not. + # Preserve synthetic provenance outside OHLC fields. to_push: np.ndarray = _prepend_tsdb_history( shm, tsdb_history, field_map=storemod.ohlc_key_map, ) + if synthetic and len(to_push): + low: float = float(to_push['time'][0]) + high: float = float(to_push['time'][-1]) + timeframe + synthetic = [ + (max(a, low), min(b, high)) + for a, b in synthetic + if a < high and b > low + ] + else: + synthetic = [] + _reservations[(mkt.fqme, timeframe)] = synthetic log.info( f'Loaded {to_push.shape} datums from storage' ) diff --git a/piker/tsp/_reserve.py b/piker/tsp/_reserve.py new file mode 100644 index 00000000..cb11bff2 --- /dev/null +++ b/piker/tsp/_reserve.py @@ -0,0 +1,192 @@ +# piker: trading gear for hackers +# Copyright (C) 2018-present Tyler Goodlet (in stewardship of pikers) + +# This program is free software: you can redistribute it and/or modify +# it under the terms of the GNU Affero General Public License as published by +# the Free Software Foundation, either version 3 of the License, or +# (at your option) any later version. + +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU Affero General Public License for more details. + +# You should have received a copy of the GNU Affero General Public License +# along with this program. If not, see . + +''' +Startup-only schedule reservations and explicit synthetic provenance. + +No published rows are moved. Synthetic rows never enter storage here. + +''' +from collections.abc import Awaitable, Callable +import math + +import numpy as np + +from piker.log import get_logger + +log = get_logger(__name__) +# Actor-local metadata, scoped by market and sampling period. +# Each range is half-open UTC and identifies synthetic SHM samples. +type SyntheticRanges = list[tuple[float, float]] +_reservations: dict[tuple[str, float], SyntheticRanges] = {} + + +def reservation_ranges( + fqme: str, + timeframe: float, +) -> list[tuple[float, float]]: + ''' + Return a copy for RPC readers; do not expose mutable ownership. + + ''' + return list(_reservations.get((fqme, timeframe), [])) + + +async def reserve_history_gaps( + frame: np.ndarray, + timeframe: float, + query_schedule: Callable[..., Awaitable[dict]], + max_rows: int, + newest_gaps: int = 32, +) -> tuple[np.ndarray, list[tuple[float, float]]]: + ''' + Reserve missing scheduled samples before publishing history. + + Keep venue closures compressed. Fully covered schedule replies + permit flat OHLC placeholders at the expected sampling cadence. + Unknown coverage and oversized reservations remain time gaps. + Inspect at most `newest_gaps` gaps in retained history, newest + first; zero disables reservations. Return the expanded frame + and half-open UTC synthetic ranges separately. + + `_history.tsdb_backfill()` calls this before prepending stored + history. It publishes the returned ranges in `_reservations`. + Local readers use `reservation_ranges()`; remote readers use + `piker.brokers.ib.feed.get_history_reservations()` on datad.ib. + Future segment repair will replace these slots by timestamp; + chart/FSP consumers still need synthetic-aware handling. + No synthetic rows are written to parquet here. + + TODO: coordinated insertion when coverage or capacity cannot + support a reservation; first pass never rewrites published rows. + + ''' + if newest_gaps < 0: + raise ValueError('newest_gaps must be nonnegative') + # A gap needs two observed margins and room for both. A zero + # request budget disables work without slicing `gaps[-0:]`. + if ( + max_rows < 2 + or + len(frame) < 2 + or + newest_gaps == 0 + ): + return frame, [] + frame = frame[-max_rows:] + times: np.ndarray = frame['time'] + # A nonpositive cadence cannot define sample positions. Duplicate + # or decreasing timestamps make the left/right margins ambiguous. + # Leave the retained data unchanged; this helper does not sort or + # deduplicate provider observations to guess a valid sample grid. + if ( + timeframe <= 0 + or + np.any(np.diff(times) <= 0) + ): + return frame, [] + gaps: np.ndarray = np.flatnonzero(np.diff(times) > timeframe) + inserts: dict[int, np.ndarray] = {} + ranges: list[tuple[float, float]] = [] + budget: int = max_rows + # Match IB Client.history_schedule()'s first-pass 31-day bound; + # 86400 converts elapsed days to seconds, not local venue days. + # TODO: advertise query limits through backend history config. + max_schedule_span_s: int = 31 * 86400 + for offset in reversed(gaps[-newest_gaps:]): + left: float = float(times[offset]) + right: float = float(times[offset + 1]) + if right - left > max_schedule_span_s: + continue + try: + # Backend-owned callback supplied in history config. + # IB's open_history_client() routes it through the proxy + # to Client.history_schedule(), returning UTC intervals. + schedule: dict = await query_schedule(left, right) + except Exception as exc: + log.warning(f'Schedule unavailable for {left}:{right}: ' + f'{exc!r}') + continue + if ( + schedule['start'] > left + or + schedule['end'] < right + or + schedule['use_rth'] + or + not schedule['sessions'] + ): + continue + stamps: list[np.ndarray] = [] + count: int = 0 + for start, end in schedule['sessions']: + first: float = max(left + timeframe, start) + # Preserve the existing sample grid across closures. + first = left + math.ceil( + (first - left) / timeframe + ) * timeframe + stop: float = min(right, end) + size: int = max(0, math.ceil((stop - first) / timeframe)) + count += size + if count > budget: + break + if size: + stamps.append(first + np.arange(size) * timeframe) + if not count: + continue + if count > budget: + # Recoverable: leave this entire gap packed/unresolved. + # No rows or provenance have been published for it, so + # later gaps can still use the remaining reservation + # budget. Do not partially reserve or reindex this gap. + log.warning(f'Reservation exceeds SHM capacity: ' + f'{left}:{right}') + continue + missing: np.ndarray = np.concatenate(stamps) + reserved: np.ndarray = np.zeros(count, dtype=frame.dtype) + reserved['time'] = missing + # COPY CLOSE / SYNTHESIZE TIMESTAMPS for scheduled slots. + # Return explicit provenance; zero volume is NOT a flag. + reserved[['open', 'high', 'low', 'close']] = ( + frame['close'][offset] + ) + inserts[int(offset)] = reserved + budget -= count + # `np.diff` locates breaks in the synthetic sample cadence. + # Add 1 to each break's left index to select the first row + # of the next run. `np.split` returns one array per run; + # encode each as [first timestamp, last + timeframe). + # E.g. [101, 102, 108] -> [101, 103), [108, 109) at 1s. + # Separate runs keep venue closures OUT of synthetic ranges. + for group in np.split( + missing, + np.flatnonzero(np.diff(missing) != timeframe) + 1, + ): + ranges.append(( + float(group[0]), float(group[-1] + timeframe), + )) + if not inserts: + return frame, [] + pieces: list[np.ndarray] = [] + previous: int = 0 + for offset, reserved in sorted(inserts.items()): + pieces.extend([frame[previous:offset + 1], reserved]) + previous = offset + 1 + pieces.append(frame[previous:]) + result: np.ndarray = np.concatenate(pieces)[-max_rows:] + first_t: float = float(result['time'][0]) + ranges = [(max(a, first_t), b) for a, b in ranges if b > first_t] + return result, sorted(ranges) diff --git a/tests/providers/test_history_reservations.py b/tests/providers/test_history_reservations.py new file mode 100644 index 00000000..8888f089 --- /dev/null +++ b/tests/providers/test_history_reservations.py @@ -0,0 +1,248 @@ +# piker: trading gear for hackers +# Copyright (C) 2018-present Tyler Goodlet (in stewardship of pikers) + +# This program is free software: you can redistribute it and/or modify +# it under the terms of the GNU Affero General Public License as published by +# the Free Software Foundation, either version 3 of the License, or +# (at your option) any later version. + +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU Affero General Public License for more details. + +# You should have received a copy of the GNU Affero General Public License +# along with this program. If not, see . + +''' +Schedule-driven startup allocation without published-row movement. + +''' +from functools import partial +from pathlib import Path + +import numpy as np +import pytest +import trio + +from piker.data._source import def_iohlcv_fields +from piker.tsp._reserve import reserve_history_gaps + + +@pytest.mark.parametrize('case', [ + 'covered', # Reserve open-session slots; exclude the closure. + 'partial', # Coverage misses the gap's left margin. + 'timeout', # Provider failure is not evidence of a closure. + 'rth', # Regular-hours sessions omit extended trading hours. + 'capacity', # Too many placeholders: leave the gap unresolved. + 'empty', # No sessions: do not assume the venue was closed. +]) +def test_reservation_requires_covered_extended_sessions(case): + ''' + Interrupted backfill leaves packed rows around missing time. + + Model a closure inside that interval: reserve only open-session + samples and retain explicit synthetic ranges. Partial coverage, + provider failure, RTH-only replies, and insufficient capacity + must never invent a schedule. Input rows remain immutable and + output placeholders carry zero activity, not invented trades. + + ''' + frame = np.zeros(2, dtype=np.dtype(def_iohlcv_fields)) + frame['time'] = [100, 110] + frame['close'] = [7, 9] + frame['volume'] = [4, 5] + before = frame.copy() + async def query(left, right): + assert (left, right) == (100, 110) + if case == 'timeout': + raise TimeoutError + return { + 'start': 102 if case == 'partial' else 100, + 'end': 111, + 'sessions': [] if case == 'empty' else [ + (100, 104), (108, 111), + ], + 'use_rth': case == 'rth', + } + result, ranges = trio.run(partial( + reserve_history_gaps, frame, 1, query, + max_rows=3 if case == 'capacity' else 20, + )) + np.testing.assert_array_equal(frame, before) + if case == 'covered': + assert result['time'].tolist() == [100, 101, 102, 103, + 108, 109, 110] + assert ranges == [(101, 104), (108, 110)] + assert result['volume'].tolist() == [4, 0, 0, 0, 0, 0, 5] + assert result['close'].tolist() == [7, 7, 7, 7, 7, 7, 9] + else: + np.testing.assert_array_equal(result, before) + assert ranges == [] + + +@pytest.mark.parametrize('newest_gaps', [-1, 0, 1, 2]) +def test_reservation_limits_newest_gap_requests( + newest_gaps: int, +) -> None: + ''' + A configurable request budget must bound provider calls. + + Two gaps distinguish newest-first selection from array order. + Zero must disable queries, not select all gaps via `[-0:]`; + negative budgets are rejected before any provider call. Check + both fetched intervals and resulting synthetic timestamps. + + ''' + frame: np.ndarray = np.zeros(3, dtype=def_iohlcv_fields) + frame['time'] = [100, 103, 106] + calls: list[tuple[float, float]] = [] + + async def query(left: float, right: float) -> dict: + calls.append((left, right)) + return { + 'start': left, + 'end': right, + 'use_rth': False, + 'sessions': [(left, right)], + } + + run = partial( + reserve_history_gaps, + frame, 1, query, max_rows=20, newest_gaps=newest_gaps, + ) + if newest_gaps < 0: + with pytest.raises(ValueError, match='nonnegative'): + trio.run(run) + assert calls == [] + return + + result, ranges = trio.run(run) + assert calls == [(103, 106), (100, 103)][:newest_gaps] + expected: dict[int, list[int]] = { + 0: [100, 103, 106], + 1: [100, 103, 104, 105, 106], + 2: [100, 101, 102, 103, 104, 105, 106], + } + assert result['time'].tolist() == expected[newest_gaps] + assert len(ranges) == newest_gaps + np.testing.assert_array_equal(frame['time'], [100, 103, 106]) + + +def test_startup_reserves_shm_without_persisting_synthetic_rows( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + ''' + Restart after interrupted backfill must retain gap provenance. + + Seed parquet with two separated observations, then run the real + startup backfiller and SHM prepend. A fake schedule supplies a + maintenance closure inside the hole. Wait for the actual tail + repair to complete before cancelling the service task. Assert + reserved open-session rows exist in SHM with queryable metadata, + while a fresh parquet load contains only provider observations. + + ''' + from contextlib import asynccontextmanager + from types import SimpleNamespace + from uuid import uuid4 + + from pendulum import duration, from_timestamp + import tractor + + from piker.data._sharedmem import maybe_open_shm_array + from piker.storage.nativedb import ( + NativeStorageClient, ohlc_key_map, + ) + from piker.tsp import _history + from piker.tsp._reserve import ( + reservation_ranges, + ) + + epoch = 1790193489 + fqme = 'schedule.test' + finished = trio.Event() + def bars(offsets): + frame = np.zeros(len(offsets), dtype=def_iohlcv_fields) + frame['time'] = np.asarray(offsets) + epoch + for field in ('open', 'high', 'low', 'close'): + frame[field] = 100 + frame['volume'] = 1 + return frame + async def history(tf, end_dt=None, **kw): + frame = bars([11, 12] if end_dt is None else [10, 11]) + return (frame, from_timestamp(frame['time'][0]), + from_timestamp(frame['time'][-1])) + async def schedule(left, right): + return { + 'start': epoch, 'end': epoch + 13, 'use_rth': False, + 'sessions': [ + (epoch, epoch + 4), (epoch + 8, epoch + 13), + ], + } + @asynccontextmanager + async def history_client(*args, **kwargs): + yield history, { + 'frame_types': {1: duration(seconds=2)}, + 'query_schedule': schedule, + } + real_repair = _history.maybe_fill_null_segments + async def repair(**kwargs): + await real_repair(**kwargs) + finished.set() + monkeypatch.setattr(_history, 'maybe_fill_null_segments', repair) + # Isolate actor-local provenance as well as SHM and parquet. + monkeypatch.setattr( + 'piker.tsp._reserve._reservations', {}, + ) + class Sampler: + async def send(self, msg): + pass + async def main(): + store = NativeStorageClient(tmp_path) + await store.update_ohlcv(fqme, bars([0, 10]), 1) + with trio.fail_after(5): + async with tractor.open_root_actor( + name='schedule-test', + tpt_bind_addrs=[('127.0.0.1', 0)], + ): + shm, opened = maybe_open_shm_array( + key=f'schedule_{uuid4().hex}', size=100, + dtype=np.dtype(def_iohlcv_fields), + append_start_index=80, + ) + assert opened + async with trio.open_nursery() as nursery: + nursery.start_soon(partial( + _history.tsdb_backfill, + mod=SimpleNamespace( + name='fake', + open_history_client=history_client, + ), + storemod=SimpleNamespace( + ohlc_key_map=ohlc_key_map, + ), + storage=store, + mkt=SimpleNamespace( + fqme=fqme, get_fqme=lambda **kw: fqme, + dst=SimpleNamespace(atype='future'), + src=SimpleNamespace(atype='fiat'), + ), + shm=shm, timeframe=1, + sampler_stream=Sampler(), + )) + await finished.wait() + nursery.cancel_scope.cancel() + assert (shm.array['time'] - epoch).tolist() == [ + 0, 1, 2, 3, 8, 9, 10, 11, 12, + ] + assert reservation_ranges(fqme, 1) == [ + (epoch + 1, epoch + 4), (epoch + 8, epoch + 10), + ] + fresh_store = NativeStorageClient(tmp_path) + disk, _, _ = await fresh_store.load(fqme, 1) + assert (disk['time'] - epoch).tolist() == [ + 0, 10, 11, 12, + ] + trio.run(main) diff --git a/tests/providers/test_ib_schedule.py b/tests/providers/test_ib_schedule.py new file mode 100644 index 00000000..6e7615c6 --- /dev/null +++ b/tests/providers/test_ib_schedule.py @@ -0,0 +1,175 @@ +# piker: trading gear for hackers +# Copyright (C) 2018-present Tyler Goodlet (in stewardship of pikers) + +# This program is free software: you can redistribute it and/or modify +# it under the terms of the GNU Affero General Public License as published by +# the Free Software Foundation, either version 3 of the License, or +# (at your option) any later version. + +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU Affero General Public License for more details. + +# You should have received a copy of the GNU Affero General Public License +# along with this program. If not, see . + +''' +Historical schedule normalization and request-boundary regressions. + +''' +import asyncio +from datetime import UTC, datetime +from types import SimpleNamespace +from zoneinfo import ZoneInfo + +import pytest + +from piker.brokers.ib.venues import normalize_schedule +from piker.brokers.ib.api import Client + + +def response(zone='America/New_York'): + ''' + Describe two sessions separated by a maintenance closure. + + ''' + return SimpleNamespace( + timeZone=zone, + startDateTime='20260923-00:00:00', + endDateTime='20260925-00:00:00', + sessions=[SimpleNamespace( + startDateTime='20260923-18:00:00', + endDateTime='20260924-17:00:00', + )], + ) + + +def test_schedule_uses_dated_timezone_and_coverage(): + ''' + Historical allocation must use dated sessions, not today's hours. + + A New York September session must convert using EDT and preserve + the response's wider coverage separately from its open interval. + + ''' + out = normalize_schedule(response(), False) + expected = datetime(2026, 9, 23, 22, tzinfo=UTC).timestamp() + assert out['sessions'][0][0] == expected + assert out['start'] < expected < out['end'] + assert out['use_rth'] is False + + +@pytest.mark.parametrize('stamp', [ + '20261101-01:30:00', '20260308-02:30:00', +]) +def test_schedule_rejects_ambiguous_or_nonexistent_time(stamp): + ''' + A guessed DST fold can shift reservations by an hour. + + Reject both the repeated autumn hour and nonexistent spring + hour instead of turning unknown coverage into a false closure. + + ''' + schedule = response() + schedule.startDateTime = stamp + with pytest.raises(ValueError, match='Ambiguous'): + normalize_schedule(schedule, False) + + +def test_schedule_request_uses_extended_hours_and_utc(): + ''' + Schedule requests must match the bars endpoint's extended hours. + + Exercise the real Client method with a fake IB transport. Assert + exact UTC request bounds and day rounding, and ensure response + coverage is not replaced by the caller's requested interval. + + ''' + calls = [] + async def find_contracts(fqme): + assert fqme == 'mnq.cme.20261218' + return ['contract'] + async def schedule(contract, **kwargs): + calls.append((contract, kwargs)) + return response() + client = SimpleNamespace( + find_contracts=find_contracts, + ib=SimpleNamespace(reqHistoricalScheduleAsync=schedule), + ) + out = asyncio.run(Client.history_schedule( + client, 'mnq.cme.20261218', + datetime(2026, 9, 23, 0, tzinfo=UTC), + datetime(2026, 9, 24, 1, tzinfo=UTC), + )) + assert calls == [('contract', { + 'numDays': 2, 'endDateTime': '20260924-01:00:00', + 'useRTH': False, + })] + assert out['timezone'] == 'America/New_York' + + +def test_schedule_timeout_is_not_an_empty_session_list(): + ''' + A failed schedule request must not classify the range as closed. + + Block the simulated provider and apply an immediate timeout; + the adapter must propagate failure instead of returning coverage. + + ''' + async def find_contracts(fqme): + return ['contract'] + async def schedule(*args, **kwargs): + await asyncio.sleep(100) + client = SimpleNamespace( + find_contracts=find_contracts, + ib=SimpleNamespace(reqHistoricalScheduleAsync=schedule), + ) + with pytest.raises(TimeoutError): + asyncio.run(Client.history_schedule( + client, 'mnq.cme', + datetime(2026, 9, 23, tzinfo=UTC), + datetime(2026, 9, 24, tzinfo=UTC), + timeout_s=0, + )) + + +def test_schedule_rejects_reversed_session(): + ''' + A malformed session must not silently become a venue closure. + + Reject the response instead of dropping its reversed interval + and reserving slots using the remaining incomplete session list. + + ''' + schedule = response() + schedule.sessions[0].endDateTime = '20260923-17:00:00' + with pytest.raises(ValueError, match='Invalid schedule session'): + normalize_schedule(schedule, False) + + +def test_schedule_duration_counts_elapsed_time_across_dst(): + ''' + The autumn clock change makes this local day last 25 hours. + + Request two elapsed days so duration rounding cannot leave an + hour uncovered by subtracting same-zone wall-clock datetimes. + + ''' + async def find_contracts(fqme): + return ['contract'] + + async def schedule(contract, **kwargs): + assert kwargs['numDays'] == 2 + return response() + + client = SimpleNamespace( + find_contracts=find_contracts, + ib=SimpleNamespace(reqHistoricalScheduleAsync=schedule), + ) + zone = ZoneInfo('America/New_York') + asyncio.run(Client.history_schedule( + client, 'mnq.cme', + datetime(2026, 11, 1, tzinfo=zone), + datetime(2026, 11, 2, tzinfo=zone), + ))