Preserve IB history request bounds on no data
Complete HMDS no-data replies without silently retrying a day older. Let `open_history_client()` classify the original request and leave frame stepping to `start_backfill()`. Distinguish `request_end_dt` from `oldest_shm_dt` so the first frame after an empty interval retains its unwritten end bar. Inline the storage time cutoff and drop the unused two-direction `diff_history()` helper. Explain both dates and filtering rules. Add request-boundary tests and real SHM/NativeDB regressions: all 2,004 supplied bars survive, while the closure stays empty. Prompt-IO: ai/prompt-io/codex/20260923T164656Z_efa72980_prompt_io.md (this patch was generated in some part by `codex` using `gpt-6` (`openai`))wkt/ib_history_request_bounds
parent
efa72980f8
commit
1da9b760a3
|
|
@ -0,0 +1,64 @@
|
||||||
|
---
|
||||||
|
model: gpt-6 (exact variant not exposed)
|
||||||
|
service: codex
|
||||||
|
timestamp: 2026-09-23T16:46:56.830897+00:00
|
||||||
|
git_ref: efa72980
|
||||||
|
session: 01a0a5b9-f872-7ca3-ab8a-d1cffc14e895
|
||||||
|
scope: code
|
||||||
|
substantive: true
|
||||||
|
raw_file: 20260923T164656Z_efa72980_prompt_io.raw.md
|
||||||
|
---
|
||||||
|
|
||||||
|
## Prompt
|
||||||
|
|
||||||
|
Investigate existing branch fixes before repairing IB day-skipping.
|
||||||
|
Create an isolated worktree on ib_connectivity_harden and add a test
|
||||||
|
that checks actual persisted data, leaving live gap repair for later.
|
||||||
|
|
||||||
|
## Response summary
|
||||||
|
|
||||||
|
Preserve IB request boundaries on HMDS no-data responses. Retain the
|
||||||
|
existing earliest-history classification and bounded reset handling.
|
||||||
|
Separate requested and published boundaries during backfill so the
|
||||||
|
first response after an empty interval does not lose its end bar.
|
||||||
|
Prove supplied bars survive both SHM publication and a fresh parquet
|
||||||
|
load while a simulated venue closure remains empty.
|
||||||
|
|
||||||
|
## Files changed
|
||||||
|
|
||||||
|
- `piker/brokers/ib/feed.py`: complete HMDS no-data requests unchanged.
|
||||||
|
- `piker/tsp/_history.py`: filter overlap against published history.
|
||||||
|
- `tests/test_ib_history.py`: six empty-response boundary cases.
|
||||||
|
- `tests/test_history_backfill.py`: two SHM/parquet integration cases.
|
||||||
|
|
||||||
|
## Human edits
|
||||||
|
|
||||||
|
The human identified the live gap, required prior-branch discovery,
|
||||||
|
selected ib_connectivity_harden as the base, and requested stronger
|
||||||
|
end-to-end evidence. They corrected overstated reproduction claims,
|
||||||
|
insisted on plain explanations, and pointed out that valid venue
|
||||||
|
closures must remain gaps. Existing data repair was explicitly deferred.
|
||||||
|
No direct human source edits have been observed; changes await review.
|
||||||
|
|
||||||
|
## Review follow-up
|
||||||
|
|
||||||
|
The human's tuicr review asked whether None was part of the backfill
|
||||||
|
protocol and requested plain explanations of the two dates and overlap
|
||||||
|
filter. Clarified the IB-internal empty-response sentinel and its
|
||||||
|
translation into NoData/DataUnavailable before reaching the backfiller.
|
||||||
|
Renamed last_start_dt to request_end_dt and published_start_dt to
|
||||||
|
oldest_shm_dt. Added a concrete empty-interval example at the filter and
|
||||||
|
explained that only successful SHM writes move oldest_shm_dt.
|
||||||
|
These changes clarify naming and documentation without changing the
|
||||||
|
previously tested request or publication behavior.
|
||||||
|
|
||||||
|
## History-filter review follow-up
|
||||||
|
|
||||||
|
The human flagged the unused append_until_dt argument and confusing
|
||||||
|
append/prepend terminology in diff_history(). The helper had one
|
||||||
|
caller, using only prepend_until_dt. Removed the helper and inlined
|
||||||
|
its inclusive lower-time cutoff beside the SHM duplicate filter.
|
||||||
|
The human's nearby TODO comment is resolved by this removal; the
|
||||||
|
replacement comment describes the retained timestamps directly.
|
||||||
|
The storage comparison remains >= backfill_until_dt; the SHM
|
||||||
|
comparison remains < oldest_shm_dt.
|
||||||
|
|
@ -0,0 +1,87 @@
|
||||||
|
---
|
||||||
|
model: gpt-6 (exact variant not exposed)
|
||||||
|
service: codex
|
||||||
|
timestamp: 2026-09-23T16:46:56.830897+00:00
|
||||||
|
git_ref: efa72980
|
||||||
|
diff_cmd: git diff efa72980 -- piker/brokers/ib/feed.py piker/tsp/_history.py tests/test_ib_history.py tests/test_history_backfill.py
|
||||||
|
---
|
||||||
|
|
||||||
|
## Prompt
|
||||||
|
|
||||||
|
The human requested the IB request-boundary repair in a new worktree
|
||||||
|
based on ib_connectivity_harden and a test exercising real persistence.
|
||||||
|
The existing MNQ parquet gap is outside this repair scope.
|
||||||
|
|
||||||
|
## Code output
|
||||||
|
|
||||||
|
> `git diff efa72980 -- piker/brokers/ib/feed.py`
|
||||||
|
|
||||||
|
> `git diff efa72980 -- piker/tsp/_history.py`
|
||||||
|
|
||||||
|
> `git diff efa72980 -- tests/test_ib_history.py`
|
||||||
|
|
||||||
|
> `git diff efa72980 -- tests/test_history_backfill.py`
|
||||||
|
|
||||||
|
get_bars() now completes HMDS no-data responses without changing the
|
||||||
|
requested date. open_history_client() retains its NoData and earliest
|
||||||
|
history classification. start_backfill() tracks the published boundary
|
||||||
|
separately from the request cursor after empty intervals.
|
||||||
|
|
||||||
|
## Unedited progress output
|
||||||
|
|
||||||
|
The new test failed on the original code as intended: the
|
||||||
|
temporary parquet contained only 4 bars instead of 2,004.
|
||||||
|
The day-skip caused 2,000 available bars to be omitted.
|
||||||
|
|
||||||
|
The test also includes a simulated venue closure that must remain
|
||||||
|
a gap. It checks that we fetch available bars without inventing
|
||||||
|
data for an empty interval.
|
||||||
|
|
||||||
|
The IB fix now preserves all 2,004 bars in parquet. The SHM
|
||||||
|
assertion exposed a second bug: after stepping over an empty
|
||||||
|
interval, the backfiller drops the last bar of the next response
|
||||||
|
as though it were already published.
|
||||||
|
|
||||||
|
I’ll separate the request date from the published-data boundary
|
||||||
|
so the same test can verify that both memory and disk retain
|
||||||
|
every available bar.
|
||||||
|
|
||||||
|
## Validation
|
||||||
|
|
||||||
|
Using the existing fsp_backfill_sync/py313 interpreter, with cwd in
|
||||||
|
ib_history_request_bounds and the Piker import checked against that cwd:
|
||||||
|
|
||||||
|
- tests/test_history_backfill.py: 11 passed in 1.92s.
|
||||||
|
- tests/test_ib_history.py: 17 passed in 0.36s.
|
||||||
|
- Ruff on all four changed Python files: all checks passed.
|
||||||
|
- git diff --check: passed.
|
||||||
|
|
||||||
|
The integration test uses simulated IB bars responses, real IB history
|
||||||
|
adapter and backfill logic, real SHM, and temporary NativeDB parquet.
|
||||||
|
It checks a fresh disk load. It does not test a live gateway, actor RPC,
|
||||||
|
or Qt rendering, and does not prove the origin of the user's live gap.
|
||||||
|
|
||||||
|
## Unedited review response
|
||||||
|
|
||||||
|
`None` is internal to the IB adapter. The backfiller receives
|
||||||
|
`NoData`, or `DataUnavailable` when the requested date has reached
|
||||||
|
the known beginning of history. I’ll make that distinction explicit.
|
||||||
|
|
||||||
|
The two dates answer different questions: “where should the next
|
||||||
|
history query end?” and “what is the oldest bar already in SHM?”
|
||||||
|
An empty reply changes only the first. I’ll rename them to
|
||||||
|
`request_end_dt` and `oldest_shm_dt` and explain the filtering with
|
||||||
|
a concrete example.
|
||||||
|
|
||||||
|
> `git diff efa72980 -- piker/brokers/ib/feed.py piker/tsp/_history.py`
|
||||||
|
|
||||||
|
## History-filter review follow-up
|
||||||
|
|
||||||
|
The human flagged the unused append_until_dt argument and confusing
|
||||||
|
append/prepend terminology in diff_history(). The helper had one
|
||||||
|
caller, using only prepend_until_dt. Removed the helper and inlined
|
||||||
|
its inclusive lower-time cutoff beside the SHM duplicate filter.
|
||||||
|
The human's nearby TODO comment is resolved by this removal; the
|
||||||
|
replacement comment describes the retained timestamps directly.
|
||||||
|
The storage comparison remains >= backfill_until_dt; the SHM
|
||||||
|
comparison remains < oldest_shm_dt.
|
||||||
|
|
@ -38,7 +38,6 @@ from async_generator import aclosing
|
||||||
import ib_async as ibis
|
import ib_async as ibis
|
||||||
import numpy as np
|
import numpy as np
|
||||||
from pendulum import (
|
from pendulum import (
|
||||||
now,
|
|
||||||
from_timestamp,
|
from_timestamp,
|
||||||
Duration,
|
Duration,
|
||||||
duration as mk_duration,
|
duration as mk_duration,
|
||||||
|
|
@ -221,7 +220,11 @@ async def open_history_client(
|
||||||
f'mean: {mean}'
|
f'mean: {mean}'
|
||||||
)
|
)
|
||||||
|
|
||||||
# could be trying to retreive bars over weekend
|
# get_bars() uses None internally for an empty reply.
|
||||||
|
# The backfiller never receives that sentinel: translate
|
||||||
|
# it into NoData, or DataUnavailable at the known start
|
||||||
|
# of the instrument's history. A venue closure can also
|
||||||
|
# yield an empty reply without exhausting all history.
|
||||||
if out is None:
|
if out is None:
|
||||||
log.error(
|
log.error(
|
||||||
f"No bars starting at {end_dt!r} !?!?"
|
f"No bars starting at {end_dt!r} !?!?"
|
||||||
|
|
@ -456,10 +459,6 @@ async def get_bars(
|
||||||
# how long before we trigger a feed reset (seconds)
|
# how long before we trigger a feed reset (seconds)
|
||||||
feed_reset_timeout: float = 3,
|
feed_reset_timeout: float = 3,
|
||||||
|
|
||||||
# how many days to subtract before giving up on further
|
|
||||||
# history queries for instrument, presuming that most don't
|
|
||||||
# not trade for a week XD
|
|
||||||
max_nodatas: int = 6,
|
|
||||||
max_failed_resets: int = 6,
|
max_failed_resets: int = 6,
|
||||||
|
|
||||||
task_status: TaskStatus[trio.CancelScope] = trio.TASK_STATUS_IGNORED,
|
task_status: TaskStatus[trio.CancelScope] = trio.TASK_STATUS_IGNORED,
|
||||||
|
|
@ -479,7 +478,6 @@ async def get_bars(
|
||||||
`.ib.api.Client` methods.
|
`.ib.api.Client` methods.
|
||||||
|
|
||||||
'''
|
'''
|
||||||
nodatas_count: int = 0
|
|
||||||
cancelled_count: int = 0
|
cancelled_count: int = 0
|
||||||
failed_resets: int = 0
|
failed_resets: int = 0
|
||||||
pacing_reset_pending: bool = False
|
pacing_reset_pending: bool = False
|
||||||
|
|
@ -496,8 +494,8 @@ async def get_bars(
|
||||||
|
|
||||||
async def query():
|
async def query():
|
||||||
|
|
||||||
nonlocal result, data_cs, end_dt
|
nonlocal result, data_cs
|
||||||
nonlocal nodatas_count, cancelled_count, failed_resets
|
nonlocal cancelled_count, failed_resets
|
||||||
nonlocal pacing_reset_pending, pacing_reset_completed
|
nonlocal pacing_reset_pending, pacing_reset_completed
|
||||||
|
|
||||||
while True:
|
while True:
|
||||||
|
|
@ -625,27 +623,21 @@ async def get_bars(
|
||||||
# the frame dt index since the upper layer may
|
# the frame dt index since the upper layer may
|
||||||
# be doing so concurrently and we don't want to
|
# be doing so concurrently and we don't want to
|
||||||
# be delivering frames that weren't asked for.
|
# be delivering frames that weren't asked for.
|
||||||
# try to decrement start point and look further back
|
|
||||||
# end_dt = end_dt.subtract(seconds=2000)
|
|
||||||
logmsg = "SUBTRACTING DAY from DT index"
|
|
||||||
if end_dt is not None:
|
|
||||||
end_dt = end_dt.subtract(days=1)
|
|
||||||
elif end_dt is None:
|
|
||||||
end_dt = now().subtract(days=1)
|
|
||||||
|
|
||||||
log.warning(
|
log.warning(
|
||||||
f'NO DATA found ending @ {end_dt}\n'
|
f'NO DATA found ending @ {end_dt}\n'
|
||||||
+ logmsg
|
|
||||||
)
|
)
|
||||||
|
# Match get_bars()'s empty-bars return: None is
|
||||||
if nodatas_count >= max_nodatas:
|
# internal to the IB adapter, not its backfill
|
||||||
raise DataUnavailable(
|
# protocol. open_history_client() converts it
|
||||||
f'Presuming {fqme} has no further history '
|
# to NoData (or DataUnavailable at history start).
|
||||||
f'after {max_nodatas} tries..'
|
# Wake the waiting get_bars() task before exiting
|
||||||
)
|
# so it does not mistake this reply for a timeout.
|
||||||
|
result = None
|
||||||
nodatas_count += 1
|
failed_resets = 0
|
||||||
continue
|
result_ready.set()
|
||||||
|
if data_cs:
|
||||||
|
data_cs.cancel()
|
||||||
|
return None
|
||||||
|
|
||||||
elif (
|
elif (
|
||||||
'API historical data query cancelled'
|
'API historical data query cancelled'
|
||||||
|
|
|
||||||
|
|
@ -145,29 +145,6 @@ _rt_buffer_start = int((_days_worth - 1) * _secs_in_day)
|
||||||
_notify_timeout_s: float = 1
|
_notify_timeout_s: float = 1
|
||||||
|
|
||||||
|
|
||||||
def diff_history(
|
|
||||||
array: np.ndarray,
|
|
||||||
append_until_dt: datetime|None = None,
|
|
||||||
prepend_until_dt: datetime|None = None,
|
|
||||||
|
|
||||||
) -> np.ndarray:
|
|
||||||
|
|
||||||
# no diffing with tsdb dt index possible..
|
|
||||||
if (
|
|
||||||
prepend_until_dt is None
|
|
||||||
and
|
|
||||||
append_until_dt is None
|
|
||||||
):
|
|
||||||
return array
|
|
||||||
|
|
||||||
times = array['time']
|
|
||||||
|
|
||||||
if append_until_dt:
|
|
||||||
return array[times < append_until_dt.timestamp()]
|
|
||||||
else:
|
|
||||||
return array[times >= prepend_until_dt.timestamp()]
|
|
||||||
|
|
||||||
|
|
||||||
def mk_storage_key(mkt: MktPair) -> str:
|
def mk_storage_key(mkt: MktPair) -> str:
|
||||||
'''
|
'''
|
||||||
Build the historical storage key for a market.
|
Build the historical storage key for a market.
|
||||||
|
|
@ -586,16 +563,24 @@ async def start_backfill(
|
||||||
# - after this loop continue to check for other gaps in the
|
# - after this loop continue to check for other gaps in the
|
||||||
# (tsdb) history and (at least report) maybe fill them
|
# (tsdb) history and (at least report) maybe fill them
|
||||||
# from new frame queries to the backend?
|
# from new frame queries to the backend?
|
||||||
last_start_dt: datetime = backfill_from_dt
|
# start_backfill() asks get_hist() for bars ending at
|
||||||
|
# request_end_dt. On NoData we move that date backward
|
||||||
|
# and try an earlier interval, even though no bars arrived.
|
||||||
|
request_end_dt: datetime = backfill_from_dt
|
||||||
|
|
||||||
|
# oldest_shm_dt is the earliest bar already written to SHM
|
||||||
|
# by this backfill. It changes only after a successful push.
|
||||||
|
# NoData changes request_end_dt but leaves this date alone.
|
||||||
|
oldest_shm_dt: datetime = backfill_from_dt
|
||||||
next_prepend_index: int = backfill_from_shm_index
|
next_prepend_index: int = backfill_from_shm_index
|
||||||
|
|
||||||
est = timezone('EST')
|
est = timezone('EST')
|
||||||
|
|
||||||
while last_start_dt > backfill_until_dt:
|
while request_end_dt > backfill_until_dt:
|
||||||
log.info(
|
log.info(
|
||||||
f'Requesting {timeframe}s frame:\n'
|
f'Requesting {timeframe}s frame:\n'
|
||||||
f'backfill_until_dt: {backfill_until_dt}\n'
|
f'backfill_until_dt: {backfill_until_dt}\n'
|
||||||
f'last_start_dt: {last_start_dt}\n'
|
f'request_end_dt: {request_end_dt}\n'
|
||||||
)
|
)
|
||||||
try:
|
try:
|
||||||
(
|
(
|
||||||
|
|
@ -604,22 +589,23 @@ async def start_backfill(
|
||||||
next_end_dt,
|
next_end_dt,
|
||||||
) = await get_hist(
|
) = await get_hist(
|
||||||
timeframe,
|
timeframe,
|
||||||
end_dt=(end_dt_param := last_start_dt),
|
end_dt=(end_dt_param := request_end_dt),
|
||||||
)
|
)
|
||||||
except NoData as nodata:
|
except NoData as nodata:
|
||||||
_nodata = nodata
|
_nodata = nodata
|
||||||
orig_last_start_dt: datetime = last_start_dt
|
previous_request_end_dt: datetime = request_end_dt
|
||||||
gap_report: str = (
|
gap_report: str = (
|
||||||
f'EMPTY FRAME for `end_dt: {last_start_dt}`?\n'
|
f'EMPTY FRAME for `end_dt: {request_end_dt}`?\n'
|
||||||
f'{mod.name} -> tf@fqme: {timeframe}@{mkt.fqme}\n'
|
f'{mod.name} -> tf@fqme: {timeframe}@{mkt.fqme}\n'
|
||||||
f'last_start_dt: {orig_last_start_dt}\n\n'
|
f'request_end_dt: {previous_request_end_dt}\n'
|
||||||
|
f'\n'
|
||||||
f'bf_until: {backfill_until_dt}\n'
|
f'bf_until: {backfill_until_dt}\n'
|
||||||
)
|
)
|
||||||
# EMPTY FRAME signal with 3 (likely) causes:
|
# EMPTY FRAME signal with 3 (likely) causes:
|
||||||
#
|
#
|
||||||
# 1. range contains legit gap in venue history
|
# 1. range contains legit gap in venue history
|
||||||
# 2. history actually (edge case) **began** at the
|
# 2. history actually (edge case) **began** at the
|
||||||
# value `last_start_dt`
|
# value `request_end_dt`
|
||||||
# 3. some other unknown error (ib blocking the
|
# 3. some other unknown error (ib blocking the
|
||||||
# history-query bc they don't want you seeing how
|
# history-query bc they don't want you seeing how
|
||||||
# they cucked all the tinas.. like with options
|
# they cucked all the tinas.. like with options
|
||||||
|
|
@ -630,13 +616,13 @@ async def start_backfill(
|
||||||
# as maybe indicated by the backend to see if we
|
# as maybe indicated by the backend to see if we
|
||||||
# can get older data before this possible
|
# can get older data before this possible
|
||||||
# "history gap".
|
# "history gap".
|
||||||
last_start_dt: datetime = last_start_dt.subtract(
|
request_end_dt = request_end_dt.subtract(
|
||||||
seconds=def_frame_duration.total_seconds()
|
seconds=def_frame_duration.total_seconds()
|
||||||
)
|
)
|
||||||
gap_report += (
|
gap_report += (
|
||||||
f'Decrementing `end_dt` and retrying with,\n'
|
f'Decrementing `end_dt` and retrying with,\n'
|
||||||
f'def_frame_duration: {def_frame_duration}\n'
|
f'def_frame_duration: {def_frame_duration}\n'
|
||||||
f'(new) last_start_dt: {last_start_dt}\n'
|
f'(new) request_end_dt: {request_end_dt}\n'
|
||||||
)
|
)
|
||||||
log.warning(gap_report)
|
log.warning(gap_report)
|
||||||
# skip writing to shm/tsdb and try the next
|
# skip writing to shm/tsdb and try the next
|
||||||
|
|
@ -657,7 +643,7 @@ async def start_backfill(
|
||||||
|
|
||||||
f'fqme: {mkt.fqme}\n'
|
f'fqme: {mkt.fqme}\n'
|
||||||
f'timeframe: {timeframe}\n'
|
f'timeframe: {timeframe}\n'
|
||||||
f'last_start_dt: {last_start_dt}\n'
|
f'request_end_dt: {request_end_dt}\n'
|
||||||
f'bf_until: {backfill_until_dt}\n'
|
f'bf_until: {backfill_until_dt}\n'
|
||||||
)
|
)
|
||||||
# UGH: what's a better way?
|
# UGH: what's a better way?
|
||||||
|
|
@ -712,7 +698,7 @@ async def start_backfill(
|
||||||
# await tractor.pause()
|
# await tractor.pause()
|
||||||
|
|
||||||
expected_dur: Interval = (
|
expected_dur: Interval = (
|
||||||
last_start_dt.subtract(
|
request_end_dt.subtract(
|
||||||
seconds=timeframe
|
seconds=timeframe
|
||||||
# ^XXX, always "up to" the bar *before*
|
# ^XXX, always "up to" the bar *before*
|
||||||
)
|
)
|
||||||
|
|
@ -747,7 +733,8 @@ async def start_backfill(
|
||||||
log.warning(
|
log.warning(
|
||||||
f'{timeframe}s-series {reason} detected!\n'
|
f'{timeframe}s-series {reason} detected!\n'
|
||||||
f'fqme: {mkt.fqme}\n'
|
f'fqme: {mkt.fqme}\n'
|
||||||
f'last_start_dt: {last_start_dt}\n\n'
|
f'request_end_dt: {request_end_dt}\n'
|
||||||
|
f'\n'
|
||||||
f'recv interval: {recv_frame_dur}\n'
|
f'recv interval: {recv_frame_dur}\n'
|
||||||
f'expected interval: {expected_dur}\n\n'
|
f'expected interval: {expected_dur}\n\n'
|
||||||
|
|
||||||
|
|
@ -756,26 +743,33 @@ async def start_backfill(
|
||||||
)
|
)
|
||||||
# await tractor.pause()
|
# await tractor.pause()
|
||||||
|
|
||||||
to_store: np.ndarray = diff_history(
|
# Stop retrieving history at backfill_until_dt. Keep
|
||||||
array,
|
# that timestamp for the storage merge, plus newer bars;
|
||||||
prepend_until_dt=backfill_until_dt,
|
# discard any older bars included in this provider frame.
|
||||||
)
|
to_store: np.ndarray = array[
|
||||||
# The requested end bar is already the oldest published
|
array['time'] >= backfill_until_dt.timestamp()
|
||||||
# sample. Keep it for storage merge, but do not
|
]
|
||||||
# consume another physical SHM row for the overlap.
|
# Keep only bars older than SHM's earliest written bar.
|
||||||
|
# A bar with the same timestamp already exists in SHM;
|
||||||
|
# writing it again would duplicate that bar.
|
||||||
|
#
|
||||||
|
# Example: SHM starts at 12:00. An empty response moves
|
||||||
|
# request_end_dt to 11:26:40, but SHM still starts at
|
||||||
|
# 12:00. If the next reply includes an 11:26:40 bar, we
|
||||||
|
# must keep it: asking for that time did not write it.
|
||||||
to_push: np.ndarray = to_store[
|
to_push: np.ndarray = to_store[
|
||||||
to_store['time'] < last_start_dt.timestamp()
|
to_store['time'] < oldest_shm_dt.timestamp()
|
||||||
]
|
]
|
||||||
ln: int = len(to_push)
|
ln: int = len(to_push)
|
||||||
if ln:
|
if ln:
|
||||||
log.info(
|
log.info(
|
||||||
f'{ln} bars for {next_start_dt} -> {last_start_dt}'
|
f'{ln} bars for {next_start_dt} -> {request_end_dt}'
|
||||||
)
|
)
|
||||||
|
|
||||||
else:
|
else:
|
||||||
log.warning(
|
log.warning(
|
||||||
f'0 BARS TO PUSH after diff!?\n'
|
f'0 BARS TO PUSH after diff!?\n'
|
||||||
f'{next_start_dt} -> {last_start_dt}'
|
f'{next_start_dt} -> {request_end_dt}'
|
||||||
f'\n'
|
f'\n'
|
||||||
f'This might mean we rxed a gap frame which starts BEFORE,\n'
|
f'This might mean we rxed a gap frame which starts BEFORE,\n'
|
||||||
f'backfill_until_dt: {backfill_until_dt}\n'
|
f'backfill_until_dt: {backfill_until_dt}\n'
|
||||||
|
|
@ -812,7 +806,7 @@ async def start_backfill(
|
||||||
):
|
):
|
||||||
log.info(
|
log.info(
|
||||||
f'Writing {len(to_store)} frame to storage:\n'
|
f'Writing {len(to_store)} frame to storage:\n'
|
||||||
f'{next_start_dt} -> {last_start_dt}'
|
f'{next_start_dt} -> {request_end_dt}'
|
||||||
)
|
)
|
||||||
await storage.update_ohlcv(
|
await storage.update_ohlcv(
|
||||||
mk_storage_key(mkt),
|
mk_storage_key(mkt),
|
||||||
|
|
@ -835,7 +829,10 @@ async def start_backfill(
|
||||||
|
|
||||||
# decrement next prepend point
|
# decrement next prepend point
|
||||||
next_prepend_index = next_prepend_index - ln
|
next_prepend_index = next_prepend_index - ln
|
||||||
last_start_dt = next_start_dt
|
request_end_dt = next_start_dt
|
||||||
|
# The push succeeded: these bars now exist in SHM.
|
||||||
|
# Use their start to exclude duplicates next time.
|
||||||
|
oldest_shm_dt = next_start_dt
|
||||||
|
|
||||||
# Stop if we've hit buffer start
|
# Stop if we've hit buffer start
|
||||||
if next_prepend_index <= 0:
|
if next_prepend_index <= 0:
|
||||||
|
|
@ -848,7 +845,8 @@ async def start_backfill(
|
||||||
except ValueError as ve:
|
except ValueError as ve:
|
||||||
_ve = ve
|
_ve = ve
|
||||||
log.error(
|
log.error(
|
||||||
f'Shm prepend OVERRUN on: {next_start_dt} -> {last_start_dt}?'
|
f'Shm prepend OVERRUN on: {next_start_dt} -> '
|
||||||
|
f'{request_end_dt}?'
|
||||||
)
|
)
|
||||||
|
|
||||||
if next_prepend_index < ln:
|
if next_prepend_index < ln:
|
||||||
|
|
@ -867,12 +865,15 @@ async def start_backfill(
|
||||||
update_start_on_prepend=update_start_on_prepend,
|
update_start_on_prepend=update_start_on_prepend,
|
||||||
)
|
)
|
||||||
next_prepend_index = next_prepend_index - ln
|
next_prepend_index = next_prepend_index - ln
|
||||||
last_start_dt = next_start_dt
|
request_end_dt = next_start_dt
|
||||||
|
# The push succeeded: these bars now exist in SHM.
|
||||||
|
# Use their start to exclude duplicates next time.
|
||||||
|
oldest_shm_dt = next_start_dt
|
||||||
shm_full = True
|
shm_full = True
|
||||||
|
|
||||||
log.info(
|
log.info(
|
||||||
f'Shm pushed {ln} frame:\n'
|
f'Shm pushed {ln} frame:\n'
|
||||||
f'{next_start_dt} -> {last_start_dt}'
|
f'{next_start_dt} -> {request_end_dt}'
|
||||||
)
|
)
|
||||||
|
|
||||||
await notify_backfill(
|
await notify_backfill(
|
||||||
|
|
|
||||||
|
|
@ -706,3 +706,160 @@ def test_null_repair_falls_back_without_debug_pause(
|
||||||
frame[field][-1],
|
frame[field][-1],
|
||||||
]
|
]
|
||||||
assert frame['volume'].tolist() == [5, 0, 0, 0, 6]
|
assert frame['volume'].tolist() == [5, 0, 0, 0, 6]
|
||||||
|
|
||||||
|
|
||||||
|
@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)
|
||||||
|
|
|
||||||
|
|
@ -586,3 +586,72 @@ def test_completed_pacing_reset_restarts_timeout_wait(
|
||||||
|
|
||||||
trio.run(main)
|
trio.run(main)
|
||||||
assert reset_types == ['data', 'connection']
|
assert reset_types == ['data', 'connection']
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize('empty_reply', ['hmds_error', 'empty_bars'])
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
'boundary', ['latest', 'after_head', 'head'],
|
||||||
|
)
|
||||||
|
def test_empty_history_preserves_request_boundary(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
empty_reply: str,
|
||||||
|
boundary: str,
|
||||||
|
) -> None:
|
||||||
|
'''
|
||||||
|
Empty IB responses retain their original request metadata.
|
||||||
|
|
||||||
|
Exercise open_history_client() and get_bars() together. The
|
||||||
|
latest request must not fall back to yesterday; a dated request
|
||||||
|
must not issue another date or become historical exhaustion
|
||||||
|
unless it has actually reached the known first timestamp.
|
||||||
|
|
||||||
|
'''
|
||||||
|
from contextlib import asynccontextmanager
|
||||||
|
from piker.brokers import NoData
|
||||||
|
|
||||||
|
head = pdatetime(2026, 9, 18, tz='UTC')
|
||||||
|
end = {
|
||||||
|
'latest': None,
|
||||||
|
'after_head': head.add(days=4),
|
||||||
|
'head': head,
|
||||||
|
}[boundary]
|
||||||
|
requests: list = []
|
||||||
|
|
||||||
|
class Proxy:
|
||||||
|
async def maybe_get_head_time(self, **kwargs):
|
||||||
|
return head
|
||||||
|
|
||||||
|
async def bars(self, **kwargs):
|
||||||
|
requests.append(kwargs['end_dt'])
|
||||||
|
if empty_reply == 'hmds_error':
|
||||||
|
raise RequestError(
|
||||||
|
1, 162, 'HMDS query returned no data',
|
||||||
|
)
|
||||||
|
return [], np.empty(0), duration(seconds=2000)
|
||||||
|
|
||||||
|
@asynccontextmanager
|
||||||
|
async def fake_data_client():
|
||||||
|
yield Proxy()
|
||||||
|
|
||||||
|
monkeypatch.setattr(
|
||||||
|
ib_feed, 'open_data_client', fake_data_client,
|
||||||
|
)
|
||||||
|
mkt = SimpleNamespace(
|
||||||
|
dst=SimpleNamespace(atype='future'),
|
||||||
|
get_bs_fqme=lambda **kwargs: 'mnq.cme.20261218',
|
||||||
|
)
|
||||||
|
|
||||||
|
async def main() -> None:
|
||||||
|
with trio.fail_after(1):
|
||||||
|
async with ib_feed.open_history_client(mkt) as pair:
|
||||||
|
get_hist, _ = pair
|
||||||
|
expected = (
|
||||||
|
DataUnavailable if boundary == 'head' else NoData
|
||||||
|
)
|
||||||
|
with pytest.raises(expected) as exc:
|
||||||
|
await get_hist(1, end_dt=end)
|
||||||
|
if expected is NoData:
|
||||||
|
assert exc.value.info['end_dt'] == end
|
||||||
|
assert requests == [end]
|
||||||
|
|
||||||
|
trio.run(main)
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue