Compare commits
No commits in common. "0846cbd4a4197e47415689cf5e9d146ba3174c32" and "ce33deb6b9f3c002a16aa3b7bb2bdadf56b3cd1c" have entirely different histories.
0846cbd4a4
...
ce33deb6b9
|
|
@ -107,8 +107,6 @@ Deterministic or local first-pass targets:
|
||||||
- `tests/test_storage_audit.py`
|
- `tests/test_storage_audit.py`
|
||||||
- `tests/test_backfill_audit_snippet.py`
|
- `tests/test_backfill_audit_snippet.py`
|
||||||
- `tests/test_ib_history.py`
|
- `tests/test_ib_history.py`
|
||||||
- `tests/test_history_backfill.py`
|
|
||||||
- `tests/test_ldshm.py`
|
|
||||||
- `tests/test_accounting.py::test_account_file_default_empty`
|
- `tests/test_accounting.py::test_account_file_default_empty`
|
||||||
- `tests/test_services.py::test_runtime_boot`
|
- `tests/test_services.py::test_runtime_boot`
|
||||||
- `tests/test_services.py::test_datad_spawn`
|
- `tests/test_services.py::test_datad_spawn`
|
||||||
|
|
@ -192,8 +190,6 @@ tests/
|
||||||
test_ems.py actor, EMS, and paper-position behavior
|
test_ems.py actor, EMS, and paper-position behavior
|
||||||
test_feeds.py live Binance/Kraken feeds and shared memory
|
test_feeds.py live Binance/Kraken feeds and shared memory
|
||||||
test_ib_history.py deterministic IB history request formatting
|
test_ib_history.py deterministic IB history request formatting
|
||||||
test_history_backfill.py deterministic history/SHM orchestration
|
|
||||||
test_ldshm.py ldshm unpublished-slot guard
|
|
||||||
test_questrade.py obsolete credentialed tests; skipped
|
test_questrade.py obsolete credentialed tests; skipped
|
||||||
test_services.py pikerd/datad/feed/EMS actor lifecycle
|
test_services.py pikerd/datad/feed/EMS actor lifecycle
|
||||||
test_storage_audit.py read-only NativeDB audit and JSON CLI
|
test_storage_audit.py read-only NativeDB audit and JSON CLI
|
||||||
|
|
@ -209,9 +205,7 @@ tests/
|
||||||
| `piker/watchlists/` | `tests/test_watchlists.py` | CLI suite is skipped |
|
| `piker/watchlists/` | `tests/test_watchlists.py` | CLI suite is skipped |
|
||||||
| `piker/storage/_audit.py`, `piker/storage/cli.py` | `tests/test_storage_audit.py` | direct Typer app, no actor |
|
| `piker/storage/_audit.py`, `piker/storage/cli.py` | `tests/test_storage_audit.py` | direct Typer app, no actor |
|
||||||
| `snippets/nativedb_backfill_audit.xsh` | `tests/test_backfill_audit_snippet.py` | disposable paths only |
|
| `snippets/nativedb_backfill_audit.xsh` | `tests/test_backfill_audit_snippet.py` | disposable paths only |
|
||||||
| `piker/brokers/ib/api.py`, `feed.py` history | `tests/test_ib_history.py` | fake client, no network |
|
| `piker/brokers/ib/api.py` history formatting | `tests/test_ib_history.py` | fake client, no network |
|
||||||
| `piker/tsp/_history.py` | `tests/test_history_backfill.py` | fake provider/storage/SHM |
|
|
||||||
| `piker/storage/cli.py` ldshm null-slot guard | `tests/test_ldshm.py` | synthetic timestamps, no SHM mutation |
|
|
||||||
| `piker/config.py` | `test_account_file_default_empty` | root-network test has a known mismatch |
|
| `piker/config.py` | `test_account_file_default_empty` | root-network test has a known mismatch |
|
||||||
| `piker/accounting/` | targeted accounting node | some cases use live/configured state |
|
| `piker/accounting/` | targeted accounting node | some cases use live/configured state |
|
||||||
| `piker/ui/_style.py`, `piker/ui/qt.py` | `tests/test_dpi_font.py` | GUI/config-isolated opt-in |
|
| `piker/ui/_style.py`, `piker/ui/qt.py` | `tests/test_dpi_font.py` | GUI/config-isolated opt-in |
|
||||||
|
|
|
||||||
|
|
@ -1,40 +0,0 @@
|
||||||
---
|
|
||||||
model: gpt-5.6-sol
|
|
||||||
provider: openai
|
|
||||||
service: opencode
|
|
||||||
session: 24fb9765-a550-4570-8350-f0fc9b7e17db
|
|
||||||
timestamp: 2026-07-29T04:17:23Z
|
|
||||||
git_ref: ce33deb6
|
|
||||||
scope: code
|
|
||||||
substantive: true
|
|
||||||
raw_file: 20260729T041723Z_ce33deb6_prompt_io.raw.md
|
|
||||||
---
|
|
||||||
|
|
||||||
## Prompt
|
|
||||||
|
|
||||||
Diagnose why the chart still showed a zero first OHLCV row after the
|
|
||||||
epoch-zero timestamp repair, causing auto-y-ranging to include zero.
|
|
||||||
|
|
||||||
## Response summary
|
|
||||||
|
|
||||||
Audited live 60s and 1s NativeDB files and found one exact 60s first-slot
|
|
||||||
sentinel at epoch 60 with a completely zero OHLCV payload. Added narrowly
|
|
||||||
scoped read and merge repair, matching incoming validation, sentinel-only
|
|
||||||
startup handling, canonical reindexing, and explicit audit classification.
|
|
||||||
Preserved modern all-zero bars as unclassified evidence rather than assuming
|
|
||||||
all zero-priced instruments are corrupt.
|
|
||||||
|
|
||||||
## Files changed
|
|
||||||
|
|
||||||
- `piker/storage/nativedb.py` - first-slot sentinel repair and validation
|
|
||||||
- `piker/storage/_audit.py` - sentinel classification and zero-row evidence
|
|
||||||
- `tests/test_storage_nativedb.py` - hydration, merge, and boundary regressions
|
|
||||||
- `tests/test_storage_audit.py` - sentinel and modern-zero audit regressions
|
|
||||||
- `ai/prompt-io/opencode/20260729T041723Z_ce33deb6_prompt_io.raw.md`
|
|
||||||
- unedited response record
|
|
||||||
- `ai/prompt-io/opencode/20260729T041723Z_ce33deb6_prompt_io.md`
|
|
||||||
- provenance metadata and response summary
|
|
||||||
|
|
||||||
## Human edits
|
|
||||||
|
|
||||||
None - generated changes have not been edited by the human.
|
|
||||||
|
|
@ -1,31 +0,0 @@
|
||||||
---
|
|
||||||
model: gpt-5.6-sol
|
|
||||||
provider: openai
|
|
||||||
service: opencode
|
|
||||||
timestamp: 2026-07-29T04:17:23Z
|
|
||||||
git_ref: ce33deb6
|
|
||||||
diff_cmd: git diff HEAD~1..HEAD
|
|
||||||
---
|
|
||||||
|
|
||||||
> `git diff HEAD~1..HEAD -- piker/storage/nativedb.py piker/storage/_audit.py tests/test_storage_nativedb.py tests/test_storage_audit.py`
|
|
||||||
|
|
||||||
Fixed the remaining legacy first-slot sentinel observed in live MNQ 60s
|
|
||||||
history: `time=60` with zero open, high, low, close, and volume. NativeDB now
|
|
||||||
repairs the exact sentinel before read hydration and merge, reindexes the
|
|
||||||
filtered view, and persists the repaired frame on the next valid update.
|
|
||||||
|
|
||||||
The repair remains narrowly scoped. Modern all-zero OHLC bars stay valid for
|
|
||||||
spreads or synthetic instruments, null-containing malformed rows still fail
|
|
||||||
validation, new writes can not recreate the first-slot sentinel, and a file
|
|
||||||
containing only sentinels loads as no history so fresh backfill can proceed.
|
|
||||||
|
|
||||||
The raw audit reports total all-zero price rows, classifies first-slot
|
|
||||||
sentinels as structural violations, and leaves later all-zero rows visible as
|
|
||||||
unclassified warnings.
|
|
||||||
|
|
||||||
Verification generated with the patch:
|
|
||||||
|
|
||||||
- NativeDB, audit, history, xonsh, and IB regressions: 61 passed
|
|
||||||
- focused Ruff: passed
|
|
||||||
- `git diff --check`: passed
|
|
||||||
- final adversarial review: no findings
|
|
||||||
|
|
@ -1,34 +0,0 @@
|
||||||
---
|
|
||||||
model: openai/gpt-5.6-sol
|
|
||||||
service: opencode
|
|
||||||
session: 6409e3bb-1a9a-4ada-b6c4-c554115143e9
|
|
||||||
timestamp: 2026-07-29T18:54:06Z
|
|
||||||
git_ref: d359b1cb
|
|
||||||
scope: code
|
|
||||||
substantive: true
|
|
||||||
raw_file: 20260729T185406Z_d359b1cb_prompt_io.raw.md
|
|
||||||
---
|
|
||||||
|
|
||||||
## Prompt
|
|
||||||
|
|
||||||
The user reported that QQQ 60-second zeros appeared as a flat line after
|
|
||||||
reload, clarified that the outstanding concern was the prior `ldshm`
|
|
||||||
failure, and instructed the agent to implement the more focused fix.
|
|
||||||
|
|
||||||
## Response summary
|
|
||||||
|
|
||||||
Changed `ldshm` to classify invalid SHM snapshots before downstream
|
|
||||||
analysis. Cadence is inferred only from observed finite positive steps;
|
|
||||||
buffers containing non-finite or non-positive timestamps, or lacking a
|
|
||||||
positive cadence, are skipped rather than partially replacement-written
|
|
||||||
or reloaded. Added deterministic helper and command-level regressions.
|
|
||||||
|
|
||||||
## Files changed
|
|
||||||
|
|
||||||
- `piker/storage/cli.py` - inspect cadence and reject invalid SHM snapshots before mutation.
|
|
||||||
- `tests/test_ldshm.py` - cover QQQ-like zeros, corrupt timestamps, cadence selection, degenerate buffers, and mutation prevention.
|
|
||||||
- `.claude/skills/run-tests/test-harness-reference.md` - register the deterministic regression target.
|
|
||||||
|
|
||||||
## Human edits
|
|
||||||
|
|
||||||
None - generated output remains uncommitted.
|
|
||||||
|
|
@ -1,58 +0,0 @@
|
||||||
---
|
|
||||||
model: openai/gpt-5.6-sol
|
|
||||||
service: opencode
|
|
||||||
timestamp: 2026-07-29T18:54:06Z
|
|
||||||
git_ref: d359b1cb
|
|
||||||
diff_cmd: git diff HEAD~1..HEAD
|
|
||||||
---
|
|
||||||
|
|
||||||
The user clarified that the QQQ 60-second chart now displayed a flat
|
|
||||||
forward-filled segment, then asked whether the previously reported
|
|
||||||
`piker store ldshm` failure had been fixed. After being told it had not,
|
|
||||||
the user instructed: "ok but i read above you have a more focussed fix
|
|
||||||
todo now? so have at it!"
|
|
||||||
|
|
||||||
The implementation focused on the exact failure boundary. It does not
|
|
||||||
filter invalid rows into a partial replacement and does not mutate live
|
|
||||||
SHM. Instead, `ldshm` inspects one captured Polars snapshot, infers its
|
|
||||||
cadence from observed finite positive timestamp steps, and skips the
|
|
||||||
entire buffer before analysis, persistence, or reload when any timestamp
|
|
||||||
is non-finite or non-positive.
|
|
||||||
|
|
||||||
> `git diff HEAD~1..HEAD -- piker/storage/cli.py`
|
|
||||||
|
|
||||||
Generated `_shm_period_and_invalid_count()` and integrated its result in
|
|
||||||
`ldshm`. The helper selects the most frequent observed positive step,
|
|
||||||
using the smallest observed step as the deterministic tie-break. Short
|
|
||||||
or cadence-degenerate buffers are skipped. Invalid snapshots are logged
|
|
||||||
and cannot reach deduplication, NativeDB replacement, or SHM reload.
|
|
||||||
|
|
||||||
> `git diff HEAD~1..HEAD -- tests/test_ldshm.py`
|
|
||||||
|
|
||||||
Generated deterministic regressions for the QQQ-shaped interior zero
|
|
||||||
run, observed-step cadence selection, all invalid timestamp classes,
|
|
||||||
short and degenerate buffers, and command-level prevention of dedupe,
|
|
||||||
storage writes, and reload. The command-level test uses an immutable
|
|
||||||
Polars snapshot and failing spies, without live SHM mutation or timing.
|
|
||||||
|
|
||||||
> `git diff HEAD~1..HEAD -- .claude/skills/run-tests/test-harness-reference.md`
|
|
||||||
|
|
||||||
Added the deterministic `ldshm` regression target and its change-to-test
|
|
||||||
mapping to the repository test harness reference.
|
|
||||||
|
|
||||||
Verification output:
|
|
||||||
|
|
||||||
```text
|
|
||||||
.......................................................... [100%]
|
|
||||||
58 passed in 0.70s
|
|
||||||
```
|
|
||||||
|
|
||||||
`git diff --check` passed. Ruff was unavailable in the worktree's
|
|
||||||
existing environment. Three adversarial review passes found and drove
|
|
||||||
fixes for partial-replacement data loss, invented median timeframes,
|
|
||||||
mixed SHM observations, and pre-guard null-analysis failures. The final
|
|
||||||
review reported no remaining findings.
|
|
||||||
|
|
||||||
The patch intentionally does not claim that `ldshm` is safe to run as a
|
|
||||||
general concurrent repair writer. It handles invalid live snapshots by
|
|
||||||
refusing mutation rather than weakening NativeDB validation.
|
|
||||||
|
|
@ -1,37 +0,0 @@
|
||||||
---
|
|
||||||
model: openai/gpt-5.6-sol
|
|
||||||
service: opencode
|
|
||||||
session: 253ade33-07ba-45b5-8839-be35d17ec164
|
|
||||||
timestamp: 2026-07-29T20:40:46Z
|
|
||||||
git_ref: ad6c560f
|
|
||||||
scope: code
|
|
||||||
substantive: true
|
|
||||||
raw_file: 20260729T204046Z_ad6c560f_prompt_io.raw.md
|
|
||||||
---
|
|
||||||
|
|
||||||
## Prompt
|
|
||||||
|
|
||||||
The user reported that the QQQ 1-second chart still contained a gap after
|
|
||||||
the `ldshm` fix, asked the agent to inspect it, and authorized continued
|
|
||||||
implementation after a short pause and shutdown of the live workspace.
|
|
||||||
|
|
||||||
## Response summary
|
|
||||||
|
|
||||||
Diagnosed a 653-row active RT SHM null reservation caused by inclusive IB
|
|
||||||
reverse-query overlaps and startup repair orchestration. Prevented extra
|
|
||||||
physical SHM overlap rows, serialized and bounded IB history recovery,
|
|
||||||
made blank/error responses terminate, and constrained synthetic repair to
|
|
||||||
valid interior timestamp boundaries. Added deterministic concurrency and
|
|
||||||
data-layout regressions.
|
|
||||||
|
|
||||||
## Files changed
|
|
||||||
|
|
||||||
- `piker/brokers/ib/feed.py` - bound and serialize IB history reset/retry behavior.
|
|
||||||
- `piker/tsp/_history.py` - prevent SHM overlap skew and safely finalize null repair.
|
|
||||||
- `tests/test_ib_history.py` - cover blank, cancellation, reset, and pacing interleavings.
|
|
||||||
- `tests/test_history_backfill.py` - cover overlap and null-repair invariants.
|
|
||||||
- `.claude/skills/run-tests/test-harness-reference.md` - register deterministic coverage.
|
|
||||||
|
|
||||||
## Human edits
|
|
||||||
|
|
||||||
None - generated output remains uncommitted.
|
|
||||||
|
|
@ -1,72 +0,0 @@
|
||||||
---
|
|
||||||
model: openai/gpt-5.6-sol
|
|
||||||
service: opencode
|
|
||||||
timestamp: 2026-07-29T20:40:46Z
|
|
||||||
git_ref: ad6c560f
|
|
||||||
diff_cmd: git diff HEAD~1..HEAD
|
|
||||||
---
|
|
||||||
|
|
||||||
The user reported that QQQ still displayed a gap on the 1-second chart
|
|
||||||
after the focused `ldshm` fix and asked the agent to inspect it. The user
|
|
||||||
then paused and resumed implementation while shutting down the live piker
|
|
||||||
trading workspace.
|
|
||||||
|
|
||||||
Read-only live inspection found the active datad generation held 653
|
|
||||||
contiguous zero-time rows in QQQ's 1-second RT SHM, while durable Parquet
|
|
||||||
contained no zero rows. The run occupied absolute indexes 168148 through
|
|
||||||
168800 between timestamps 1785348398 and 1785349051. The stale actor
|
|
||||||
generation and active 60-second history were separately identified.
|
|
||||||
|
|
||||||
> `git diff HEAD~1..HEAD -- piker/tsp/_history.py`
|
|
||||||
|
|
||||||
Generated history changes that exclude each inclusive reverse-query
|
|
||||||
endpoint already present in SHM while retaining the full provider frame
|
|
||||||
for NativeDB merge. Startup null repair now runs after reverse backfill,
|
|
||||||
uses bounded UI notifications, handles provider no-data without an
|
|
||||||
interactive pause, fills only exact interior zero groups, and refuses
|
|
||||||
synthetic timestamps unless both valid boundaries establish the expected
|
|
||||||
cadence.
|
|
||||||
|
|
||||||
> `git diff HEAD~1..HEAD -- piker/brokers/ib/feed.py`
|
|
||||||
|
|
||||||
Generated IB history reset changes that make blank responses complete,
|
|
||||||
bound repeated cancellation and reset failures per request, serialize
|
|
||||||
timeout and pacing farm resets with one actor-global Trio lock, carry
|
|
||||||
explicit reset results, and coordinate pacing handoff with queued timeout
|
|
||||||
waiters.
|
|
||||||
|
|
||||||
> `git diff HEAD~1..HEAD -- tests/test_history_backfill.py`
|
|
||||||
|
|
||||||
Generated regressions for inclusive endpoint exclusion, full NativeDB
|
|
||||||
provider deltas, strict synthetic gap boundaries, empty-provider local
|
|
||||||
fallback, bounded notifications, and preservation of valid margin rows.
|
|
||||||
|
|
||||||
> `git diff HEAD~1..HEAD -- tests/test_ib_history.py`
|
|
||||||
|
|
||||||
Generated deterministic regressions for blank IB responses, cancellation
|
|
||||||
limits, stalled reset cancellation, pacing reset failure limits,
|
|
||||||
cross-request reset serialization, timeout-to-pacing handoff, and
|
|
||||||
successful handoff timeout restart.
|
|
||||||
|
|
||||||
> `git diff HEAD~1..HEAD -- .claude/skills/run-tests/test-harness-reference.md`
|
|
||||||
|
|
||||||
Registered the deterministic history-backfill target and expanded the IB
|
|
||||||
history mapping to include `feed.py`.
|
|
||||||
|
|
||||||
Verification output:
|
|
||||||
|
|
||||||
```text
|
|
||||||
........................................................................ [ 93%]
|
|
||||||
..... [100%]
|
|
||||||
77 passed in 2.12s
|
|
||||||
```
|
|
||||||
|
|
||||||
`git diff --check` passed. Multiple adversarial review rounds identified
|
|
||||||
and drove fixes for unbounded IB empty/cancellation/reset paths, stale
|
|
||||||
reset ownership, actor-global failure poisoning, concurrent farm resets,
|
|
||||||
double-counted pacing handoffs, physical endpoint overlap, and unsafe
|
|
||||||
synthetic timestamp boundaries. Final review reported no findings.
|
|
||||||
|
|
||||||
Live validation was intentionally deferred because the user shut down the
|
|
||||||
trading workspace. A fresh datad/chart startup is required to regenerate
|
|
||||||
SHM and verify that QQQ no longer contains the 653-row null run.
|
|
||||||
|
|
@ -32,6 +32,7 @@ import time
|
||||||
from typing import (
|
from typing import (
|
||||||
Any,
|
Any,
|
||||||
Callable,
|
Callable,
|
||||||
|
TYPE_CHECKING,
|
||||||
)
|
)
|
||||||
|
|
||||||
from async_generator import aclosing
|
from async_generator import aclosing
|
||||||
|
|
@ -74,6 +75,9 @@ from ._util import (
|
||||||
)
|
)
|
||||||
from .symbols import get_mkt_info
|
from .symbols import get_mkt_info
|
||||||
|
|
||||||
|
if TYPE_CHECKING:
|
||||||
|
from trio._core._run import Task
|
||||||
|
|
||||||
log = get_logger(
|
log = get_logger(
|
||||||
name=__name__,
|
name=__name__,
|
||||||
)
|
)
|
||||||
|
|
@ -349,7 +353,6 @@ async def wait_on_data_reset(
|
||||||
tuple[
|
tuple[
|
||||||
trio.CancelScope,
|
trio.CancelScope,
|
||||||
trio.Event,
|
trio.Event,
|
||||||
list[bool],
|
|
||||||
]
|
]
|
||||||
] = trio.TASK_STATUS_IGNORED,
|
] = trio.TASK_STATUS_IGNORED,
|
||||||
) -> bool:
|
) -> bool:
|
||||||
|
|
@ -393,9 +396,8 @@ async def wait_on_data_reset(
|
||||||
client: Client = proxy._aio_ns
|
client: Client = proxy._aio_ns
|
||||||
|
|
||||||
done = trio.Event()
|
done = trio.Event()
|
||||||
result: list[bool] = []
|
|
||||||
with trio.move_on_after(timeout) as cs:
|
with trio.move_on_after(timeout) as cs:
|
||||||
task_status.started((cs, done, result))
|
task_status.started((cs, done))
|
||||||
|
|
||||||
log.warning(
|
log.warning(
|
||||||
'Sending DATA RESET request:\n'
|
'Sending DATA RESET request:\n'
|
||||||
|
|
@ -411,7 +413,6 @@ async def wait_on_data_reset(
|
||||||
'NO VNC DETECTED!\n'
|
'NO VNC DETECTED!\n'
|
||||||
'Manually press ctrl-alt-f on your IB java app'
|
'Manually press ctrl-alt-f on your IB java app'
|
||||||
)
|
)
|
||||||
result.append(False)
|
|
||||||
done.set()
|
done.set()
|
||||||
return False
|
return False
|
||||||
|
|
||||||
|
|
@ -426,7 +427,6 @@ async def wait_on_data_reset(
|
||||||
]:
|
]:
|
||||||
await ev.wait()
|
await ev.wait()
|
||||||
log.info(f"{name} DATA RESET")
|
log.info(f"{name} DATA RESET")
|
||||||
result.append(True)
|
|
||||||
done.set()
|
done.set()
|
||||||
return True
|
return True
|
||||||
|
|
||||||
|
|
@ -435,12 +435,12 @@ async def wait_on_data_reset(
|
||||||
'Data reset task canceled?'
|
'Data reset task canceled?'
|
||||||
)
|
)
|
||||||
|
|
||||||
result.append(False)
|
|
||||||
done.set()
|
done.set()
|
||||||
return False
|
return False
|
||||||
|
|
||||||
|
|
||||||
_data_reset_lock = trio.Lock()
|
_data_resetter_task: Task | None = None
|
||||||
|
_failed_resets: int = 0
|
||||||
|
|
||||||
|
|
||||||
async def get_bars(
|
async def get_bars(
|
||||||
|
|
@ -479,11 +479,8 @@ async def get_bars(
|
||||||
`.ib.api.Client` methods.
|
`.ib.api.Client` methods.
|
||||||
|
|
||||||
'''
|
'''
|
||||||
|
global _data_resetter_task, _failed_resets
|
||||||
nodatas_count: int = 0
|
nodatas_count: int = 0
|
||||||
cancelled_count: int = 0
|
|
||||||
failed_resets: int = 0
|
|
||||||
pacing_reset_pending: bool = False
|
|
||||||
pacing_reset_completed: bool = False
|
|
||||||
|
|
||||||
data_cs: trio.CancelScope | None = None
|
data_cs: trio.CancelScope | None = None
|
||||||
result: tuple[
|
result: tuple[
|
||||||
|
|
@ -496,16 +493,10 @@ async def get_bars(
|
||||||
|
|
||||||
async def query():
|
async def query():
|
||||||
|
|
||||||
nonlocal result, data_cs, end_dt
|
global _failed_resets
|
||||||
nonlocal nodatas_count, cancelled_count, failed_resets
|
nonlocal result, data_cs, end_dt, nodatas_count
|
||||||
nonlocal pacing_reset_pending, pacing_reset_completed
|
|
||||||
|
|
||||||
while True:
|
while _failed_resets < max_failed_resets:
|
||||||
if failed_resets >= max_failed_resets:
|
|
||||||
raise DataUnavailable(
|
|
||||||
f'IB history reset failed {failed_resets} '
|
|
||||||
f'times for {fqme}'
|
|
||||||
)
|
|
||||||
try:
|
try:
|
||||||
(
|
(
|
||||||
bars,
|
bars,
|
||||||
|
|
@ -537,10 +528,6 @@ async def get_bars(
|
||||||
)
|
)
|
||||||
# NOTE: REQUIRED to pass back value..
|
# NOTE: REQUIRED to pass back value..
|
||||||
result = None
|
result = None
|
||||||
failed_resets = 0
|
|
||||||
result_ready.set()
|
|
||||||
if data_cs:
|
|
||||||
data_cs.cancel()
|
|
||||||
return None
|
return None
|
||||||
|
|
||||||
# not enough bars signal, likely due to venue
|
# not enough bars signal, likely due to venue
|
||||||
|
|
@ -600,7 +587,6 @@ async def get_bars(
|
||||||
first_dt,
|
first_dt,
|
||||||
last_dt,
|
last_dt,
|
||||||
)
|
)
|
||||||
failed_resets = 0
|
|
||||||
|
|
||||||
# signal data reset loop parent task
|
# signal data reset loop parent task
|
||||||
result_ready.set()
|
result_ready.set()
|
||||||
|
|
@ -656,12 +642,6 @@ async def get_bars(
|
||||||
'Query cancelled by IB (:eyeroll:):\n'
|
'Query cancelled by IB (:eyeroll:):\n'
|
||||||
f'{err.message}'
|
f'{err.message}'
|
||||||
)
|
)
|
||||||
cancelled_count += 1
|
|
||||||
if cancelled_count >= max_failed_resets:
|
|
||||||
raise DataUnavailable(
|
|
||||||
f'IB cancelled {cancelled_count} history '
|
|
||||||
f'queries for {fqme}'
|
|
||||||
)
|
|
||||||
continue
|
continue
|
||||||
|
|
||||||
elif (
|
elif (
|
||||||
|
|
@ -682,7 +662,6 @@ async def get_bars(
|
||||||
)
|
)
|
||||||
|
|
||||||
client = proxy._aio_ns.ib.client
|
client = proxy._aio_ns.ib.client
|
||||||
pacing_reset_pending = True
|
|
||||||
|
|
||||||
# cancel any existing reset task
|
# cancel any existing reset task
|
||||||
if data_cs:
|
if data_cs:
|
||||||
|
|
@ -690,81 +669,69 @@ async def get_bars(
|
||||||
data_cs.cancel()
|
data_cs.cancel()
|
||||||
|
|
||||||
# spawn new data reset task
|
# spawn new data reset task
|
||||||
async with _data_reset_lock:
|
data_cs, reset_done = await tn.start(
|
||||||
(
|
|
||||||
data_cs,
|
|
||||||
reset_done,
|
|
||||||
reset_result,
|
|
||||||
) = await tn.start(
|
|
||||||
partial(
|
partial(
|
||||||
wait_on_data_reset,
|
wait_on_data_reset,
|
||||||
proxy,
|
proxy,
|
||||||
reset_type='connection'
|
reset_type='connection'
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
await reset_done.wait()
|
if reset_done:
|
||||||
if reset_result[0]:
|
_failed_resets = 0
|
||||||
failed_resets = 0
|
|
||||||
else:
|
else:
|
||||||
failed_resets += 1
|
_failed_resets += 1
|
||||||
pacing_reset_pending = False
|
|
||||||
pacing_reset_completed = True
|
|
||||||
continue
|
continue
|
||||||
|
|
||||||
else:
|
else:
|
||||||
raise
|
raise
|
||||||
|
|
||||||
|
# TODO: make this global across all history task/requests
|
||||||
|
# such that simultaneous symbol queries don't try data resettingn
|
||||||
|
# too fast..
|
||||||
|
unset_resetter: bool = False
|
||||||
async with (
|
async with (
|
||||||
tractor.trionics.collapse_eg(),
|
tractor.trionics.collapse_eg(),
|
||||||
trio.open_nursery() as tn
|
trio.open_nursery() as tn
|
||||||
):
|
):
|
||||||
|
|
||||||
# Run history until a result or bounded reset failure.
|
# start history request that we allow
|
||||||
|
# to run indefinitely until a result is acquired
|
||||||
tn.start_soon(query)
|
tn.start_soon(query)
|
||||||
|
|
||||||
# Serialize farm resets across every history requester.
|
# start history reset loop which waits up to the timeout
|
||||||
|
# for a result before triggering a data feed reset.
|
||||||
while not result_ready.is_set():
|
while not result_ready.is_set():
|
||||||
|
|
||||||
with trio.move_on_after(feed_reset_timeout):
|
with trio.move_on_after(feed_reset_timeout):
|
||||||
await result_ready.wait()
|
await result_ready.wait()
|
||||||
break
|
break
|
||||||
|
|
||||||
async with _data_reset_lock:
|
if _data_resetter_task:
|
||||||
if result_ready.is_set():
|
# don't double invoke the reset hack if another
|
||||||
break
|
# requester task already has it covered.
|
||||||
if pacing_reset_completed:
|
|
||||||
pacing_reset_completed = False
|
|
||||||
continue
|
continue
|
||||||
|
|
||||||
(
|
else:
|
||||||
data_cs,
|
_data_resetter_task = trio.lowlevel.current_task()
|
||||||
reset_done,
|
unset_resetter: bool = True
|
||||||
reset_result,
|
|
||||||
) = await tn.start(
|
# spawn new data reset task
|
||||||
|
data_cs, reset_done = await tn.start(
|
||||||
partial(
|
partial(
|
||||||
wait_on_data_reset,
|
wait_on_data_reset,
|
||||||
proxy,
|
proxy,
|
||||||
reset_type='data',
|
reset_type='data',
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
|
# sync wait on reset to complete
|
||||||
await reset_done.wait()
|
await reset_done.wait()
|
||||||
|
|
||||||
if (
|
_data_resetter_task = (
|
||||||
result_ready.is_set()
|
None
|
||||||
or
|
if unset_resetter
|
||||||
reset_result[0]
|
else _data_resetter_task
|
||||||
):
|
|
||||||
failed_resets = 0
|
|
||||||
elif pacing_reset_pending:
|
|
||||||
# Pacing recovery deliberately cancelled this reset
|
|
||||||
# and will consume the retry itself.
|
|
||||||
pass
|
|
||||||
else:
|
|
||||||
failed_resets += 1
|
|
||||||
if failed_resets >= max_failed_resets:
|
|
||||||
raise DataUnavailable(
|
|
||||||
f'IB history reset failed '
|
|
||||||
f'{failed_resets} times for {fqme}'
|
|
||||||
)
|
)
|
||||||
|
assert result
|
||||||
return (
|
return (
|
||||||
result,
|
result,
|
||||||
data_cs is not None,
|
data_cs is not None,
|
||||||
|
|
|
||||||
|
|
@ -30,12 +30,6 @@ _value_fields: tuple[str, ...] = (
|
||||||
'close',
|
'close',
|
||||||
'volume',
|
'volume',
|
||||||
)
|
)
|
||||||
_price_fields: tuple[str, ...] = (
|
|
||||||
'open',
|
|
||||||
'high',
|
|
||||||
'low',
|
|
||||||
'close',
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def _utc_str(timestamp: float|int|None) -> str|None:
|
def _utc_str(timestamp: float|int|None) -> str|None:
|
||||||
|
|
@ -500,36 +494,6 @@ def audit_ohlcv_frame(
|
||||||
if counts['nonfinite'] != 0:
|
if counts['nonfinite'] != 0:
|
||||||
all_values_finite = False
|
all_values_finite = False
|
||||||
|
|
||||||
all_zero_price_rows: int|None = None
|
|
||||||
epoch_sentinel_rows: int|None = None
|
|
||||||
if all(
|
|
||||||
field in df.columns
|
|
||||||
and
|
|
||||||
df[field].dtype.is_numeric()
|
|
||||||
for field in _price_fields
|
|
||||||
):
|
|
||||||
zero_prices: pl.Expr = pl.all_horizontal([
|
|
||||||
pl.col(field) == 0
|
|
||||||
for field in _price_fields
|
|
||||||
]).fill_null(False)
|
|
||||||
all_zero_price_rows = int(
|
|
||||||
df.select(zero_prices).to_series().sum()
|
|
||||||
or 0
|
|
||||||
)
|
|
||||||
if (
|
|
||||||
'time' in df.columns
|
|
||||||
and
|
|
||||||
df['time'].dtype.is_numeric()
|
|
||||||
):
|
|
||||||
epoch_sentinel_rows = int(
|
|
||||||
df.select(
|
|
||||||
(pl.col('time') <= period_s)
|
|
||||||
&
|
|
||||||
zero_prices
|
|
||||||
).to_series().sum()
|
|
||||||
or 0
|
|
||||||
)
|
|
||||||
|
|
||||||
schema_ok: bool = bool(
|
schema_ok: bool = bool(
|
||||||
not missing
|
not missing
|
||||||
and
|
and
|
||||||
|
|
@ -552,8 +516,6 @@ def audit_ohlcv_frame(
|
||||||
violations.append('canonical_dtypes')
|
violations.append('canonical_dtypes')
|
||||||
if not all_values_finite:
|
if not all_values_finite:
|
||||||
violations.append('nonfinite_ohlcv')
|
violations.append('nonfinite_ohlcv')
|
||||||
if epoch_sentinel_rows:
|
|
||||||
violations.append('epoch_ohlc_sentinel')
|
|
||||||
if timestamps.get('nonfinite') != 0:
|
if timestamps.get('nonfinite') != 0:
|
||||||
violations.append('nonfinite_timestamps')
|
violations.append('nonfinite_timestamps')
|
||||||
if (
|
if (
|
||||||
|
|
@ -578,15 +540,6 @@ def audit_ohlcv_frame(
|
||||||
warnings: list[str] = []
|
warnings: list[str] = []
|
||||||
if gaps['count']:
|
if gaps['count']:
|
||||||
warnings.append('positive_time_gaps_unclassified')
|
warnings.append('positive_time_gaps_unclassified')
|
||||||
unclassified_zero_rows: int|None = None
|
|
||||||
if all_zero_price_rows is not None:
|
|
||||||
unclassified_zero_rows = (
|
|
||||||
all_zero_price_rows
|
|
||||||
-
|
|
||||||
(epoch_sentinel_rows or 0)
|
|
||||||
)
|
|
||||||
if unclassified_zero_rows:
|
|
||||||
warnings.append('all_zero_ohlc_prices_unclassified')
|
|
||||||
|
|
||||||
structural_ok: bool = not violations
|
structural_ok: bool = not violations
|
||||||
gap_free: bool|None = (
|
gap_free: bool|None = (
|
||||||
|
|
@ -620,9 +573,6 @@ def audit_ohlcv_frame(
|
||||||
'nan_by_column': nans,
|
'nan_by_column': nans,
|
||||||
'infinity_by_column': infinities,
|
'infinity_by_column': infinities,
|
||||||
'all_finite': all_values_finite,
|
'all_finite': all_values_finite,
|
||||||
'all_zero_price_rows': all_zero_price_rows,
|
|
||||||
'epoch_sentinel_rows': epoch_sentinel_rows,
|
|
||||||
'unclassified_zero_price_rows': unclassified_zero_rows,
|
|
||||||
},
|
},
|
||||||
'gaps': gaps,
|
'gaps': gaps,
|
||||||
'result': {
|
'result': {
|
||||||
|
|
|
||||||
|
|
@ -60,37 +60,6 @@ if TYPE_CHECKING:
|
||||||
store = typer.Typer()
|
store = typer.Typer()
|
||||||
|
|
||||||
|
|
||||||
def _shm_period_and_invalid_count(
|
|
||||||
times: np.ndarray,
|
|
||||||
|
|
||||||
) -> tuple[float|None, int]:
|
|
||||||
'''
|
|
||||||
Infer cadence without treating unpublished slots as samples.
|
|
||||||
|
|
||||||
'''
|
|
||||||
valid: np.ndarray = (
|
|
||||||
np.isfinite(times)
|
|
||||||
&
|
|
||||||
(times > 0)
|
|
||||||
)
|
|
||||||
invalid_count: int = int(np.count_nonzero(~valid))
|
|
||||||
published: np.ndarray = times[valid]
|
|
||||||
if published.size < 2:
|
|
||||||
return None, invalid_count
|
|
||||||
|
|
||||||
steps: np.ndarray = np.diff(published)
|
|
||||||
positive_steps: np.ndarray = steps[steps > 0]
|
|
||||||
if not positive_steps.size:
|
|
||||||
return None, invalid_count
|
|
||||||
|
|
||||||
periods, counts = np.unique(
|
|
||||||
positive_steps,
|
|
||||||
return_counts=True,
|
|
||||||
)
|
|
||||||
period: float = float(periods[np.argmax(counts)])
|
|
||||||
return period, invalid_count
|
|
||||||
|
|
||||||
|
|
||||||
def _render_audit_report(report: dict) -> None:
|
def _render_audit_report(report: dict) -> None:
|
||||||
'''
|
'''
|
||||||
Render a compact human summary and explicit gap endpoints.
|
Render a compact human summary and explicit gap endpoints.
|
||||||
|
|
@ -537,25 +506,38 @@ def ldshm(
|
||||||
shm_df,
|
shm_df,
|
||||||
) in tsp.iter_dfs_from_shms(fqme):
|
) in tsp.iter_dfs_from_shms(fqme):
|
||||||
|
|
||||||
times: np.ndarray = shm_df['time'].to_numpy()
|
times: np.ndarray = shm.array['time']
|
||||||
(
|
d1: float = float(times[-1] - times[-2])
|
||||||
period_s,
|
d2: float = 0
|
||||||
invalid_count,
|
# XXX, take a median sample rate if sufficient data
|
||||||
) = _shm_period_and_invalid_count(times)
|
if times.size > 2:
|
||||||
if period_s is None:
|
d2: float = float(times[-2] - times[-3])
|
||||||
log.warning(
|
med: float = np.median(np.diff(times))
|
||||||
f'Could not infer a positive sample period '
|
if (
|
||||||
f'for {shmfile.name}; skipping buffer\n'
|
d1 < 1.
|
||||||
|
and d2 < 1.
|
||||||
|
and med < 1.
|
||||||
|
):
|
||||||
|
raise ValueError(
|
||||||
|
f'Something is wrong with time period for {shm}:\n{times}'
|
||||||
)
|
)
|
||||||
continue
|
period_s: float = float(max(d1, d2, med))
|
||||||
log.info(
|
log.info(
|
||||||
f'Processing shm buffer:\n'
|
f'Processing shm buffer:\n'
|
||||||
f' file: {shmfile.name}\n'
|
f' file: {shmfile.name}\n'
|
||||||
f' period: {period_s}s\n'
|
f' period: {period_s}s\n'
|
||||||
)
|
)
|
||||||
|
|
||||||
|
null_segs: tuple = tsp.get_null_segs(
|
||||||
|
frame=shm.array,
|
||||||
|
period=period_s,
|
||||||
|
)
|
||||||
|
|
||||||
# TODO: call null-seg fixer somehow?
|
# TODO: call null-seg fixer somehow?
|
||||||
if invalid_count:
|
if null_segs:
|
||||||
|
|
||||||
|
if tractor.runtime._state.is_debug_mode():
|
||||||
|
await tractor.pause()
|
||||||
# async with (
|
# async with (
|
||||||
# trio.open_nursery() as tn,
|
# trio.open_nursery() as tn,
|
||||||
# mod.open_history_client(
|
# mod.open_history_client(
|
||||||
|
|
@ -571,14 +553,6 @@ def ldshm(
|
||||||
# sampler_stream=sampler_stream,
|
# sampler_stream=sampler_stream,
|
||||||
# mkt=mkt,
|
# mkt=mkt,
|
||||||
# ))
|
# ))
|
||||||
if tractor.runtime._state.is_debug_mode():
|
|
||||||
await tractor.pause()
|
|
||||||
log.warning(
|
|
||||||
f'Found {invalid_count} invalid SHM row(s) '
|
|
||||||
f'in {shmfile.name}; skipping persistence '
|
|
||||||
f'and reload\n'
|
|
||||||
)
|
|
||||||
continue
|
|
||||||
|
|
||||||
# over-write back to shm?
|
# over-write back to shm?
|
||||||
wdts: pl.DataFrame # with dts
|
wdts: pl.DataFrame # with dts
|
||||||
|
|
|
||||||
|
|
@ -74,12 +74,6 @@ from . import TimeseriesNotFound
|
||||||
|
|
||||||
|
|
||||||
log = get_logger('storage.nativedb')
|
log = get_logger('storage.nativedb')
|
||||||
_price_fields: tuple[str, ...] = (
|
|
||||||
'open',
|
|
||||||
'high',
|
|
||||||
'low',
|
|
||||||
'close',
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def detect_period(shm: ShmArray) -> float:
|
def detect_period(shm: ShmArray) -> float:
|
||||||
|
|
@ -239,8 +233,6 @@ class NativeStorageClient:
|
||||||
) from fnfe
|
) from fnfe
|
||||||
|
|
||||||
times = array['time']
|
times = array['time']
|
||||||
if array.size == 0:
|
|
||||||
return None
|
|
||||||
return (
|
return (
|
||||||
array,
|
array,
|
||||||
from_timestamp(times[0]),
|
from_timestamp(times[0]),
|
||||||
|
|
@ -327,41 +319,6 @@ class NativeStorageClient:
|
||||||
.cast(schema)
|
.cast(schema)
|
||||||
)
|
)
|
||||||
|
|
||||||
def _repair_stored_ohlcv(
|
|
||||||
self,
|
|
||||||
df: pl.DataFrame,
|
|
||||||
timeframe: int,
|
|
||||||
|
|
||||||
) -> pl.DataFrame:
|
|
||||||
'''
|
|
||||||
Drop known legacy sentinels before hydration or merge.
|
|
||||||
|
|
||||||
'''
|
|
||||||
stored: pl.DataFrame = self._canonicalize_ohlcv(df)
|
|
||||||
zero_prices: pl.Expr = pl.all_horizontal([
|
|
||||||
pl.col(field) == 0
|
|
||||||
for field in _price_fields
|
|
||||||
]).fill_null(False)
|
|
||||||
invalid: pl.Expr = (
|
|
||||||
(pl.col('time') <= 0)
|
|
||||||
|
|
|
||||||
(
|
|
||||||
(pl.col('time') <= timeframe)
|
|
||||||
&
|
|
||||||
zero_prices
|
|
||||||
)
|
|
||||||
)
|
|
||||||
stored_len: int = stored.height
|
|
||||||
stored = stored.filter(~invalid)
|
|
||||||
dropped: int = stored_len - stored.height
|
|
||||||
if dropped:
|
|
||||||
log.warning(
|
|
||||||
f'Dropping {dropped} invalid persisted OHLCV '
|
|
||||||
f'row(s) during read repair'
|
|
||||||
)
|
|
||||||
stored = self._canonicalize_ohlcv(stored)
|
|
||||||
return stored
|
|
||||||
|
|
||||||
async def read_ohlcv(
|
async def read_ohlcv(
|
||||||
self,
|
self,
|
||||||
fqme: str,
|
fqme: str,
|
||||||
|
|
@ -374,10 +331,7 @@ class NativeStorageClient:
|
||||||
fqme,
|
fqme,
|
||||||
period=int(timeframe),
|
period=int(timeframe),
|
||||||
)
|
)
|
||||||
df: pl.DataFrame = self._repair_stored_ohlcv(
|
df: pl.DataFrame = pl.read_parquet(path)
|
||||||
pl.read_parquet(path),
|
|
||||||
timeframe=int(timeframe),
|
|
||||||
)
|
|
||||||
|
|
||||||
self._cache_df(
|
self._cache_df(
|
||||||
fqme=fqme,
|
fqme=fqme,
|
||||||
|
|
@ -430,7 +384,7 @@ class NativeStorageClient:
|
||||||
else:
|
else:
|
||||||
df = ohlcv
|
df = ohlcv
|
||||||
df = self._canonicalize_ohlcv(df)
|
df = self._canonicalize_ohlcv(df)
|
||||||
self._validate_ohlcv(df, timeframe)
|
self._validate_ohlcv(df)
|
||||||
|
|
||||||
# TODO: in terms of managing the ultra long term data
|
# TODO: in terms of managing the ultra long term data
|
||||||
# -[ ] use a proper profiler to measure all this IO and
|
# -[ ] use a proper profiler to measure all this IO and
|
||||||
|
|
@ -451,7 +405,7 @@ class NativeStorageClient:
|
||||||
try:
|
try:
|
||||||
df.write_parquet(tmp_path)
|
df.write_parquet(tmp_path)
|
||||||
committed: pl.DataFrame = pl.read_parquet(tmp_path)
|
committed: pl.DataFrame = pl.read_parquet(tmp_path)
|
||||||
self._validate_ohlcv(committed, timeframe)
|
self._validate_ohlcv(committed)
|
||||||
if not committed.equals(df):
|
if not committed.equals(df):
|
||||||
raise IOError(
|
raise IOError(
|
||||||
'Temporary parquet differs from input frame'
|
'Temporary parquet differs from input frame'
|
||||||
|
|
@ -488,7 +442,6 @@ class NativeStorageClient:
|
||||||
def _validate_ohlcv(
|
def _validate_ohlcv(
|
||||||
self,
|
self,
|
||||||
df: pl.DataFrame,
|
df: pl.DataFrame,
|
||||||
timeframe: int,
|
|
||||||
|
|
||||||
) -> None:
|
) -> None:
|
||||||
'''
|
'''
|
||||||
|
|
@ -530,17 +483,6 @@ class NativeStorageClient:
|
||||||
'OHLCV timestamps must be finite and positive'
|
'OHLCV timestamps must be finite and positive'
|
||||||
)
|
)
|
||||||
|
|
||||||
prices: np.ndarray = df.select(_price_fields).to_numpy()
|
|
||||||
epoch_sentinel: np.ndarray = (
|
|
||||||
(times <= timeframe)
|
|
||||||
&
|
|
||||||
np.all(prices == 0, axis=1)
|
|
||||||
)
|
|
||||||
if np.any(epoch_sentinel):
|
|
||||||
raise ValueError(
|
|
||||||
'OHLCV frame contains a zero-initialized epoch sentinel'
|
|
||||||
)
|
|
||||||
|
|
||||||
if np.any(np.diff(times) <= 0):
|
if np.any(np.diff(times) <= 0):
|
||||||
raise ValueError(
|
raise ValueError(
|
||||||
'OHLCV timestamps must be strictly increasing'
|
'OHLCV timestamps must be strictly increasing'
|
||||||
|
|
@ -565,13 +507,20 @@ class NativeStorageClient:
|
||||||
else:
|
else:
|
||||||
incoming = ohlcv
|
incoming = ohlcv
|
||||||
incoming = self._canonicalize_ohlcv(incoming)
|
incoming = self._canonicalize_ohlcv(incoming)
|
||||||
self._validate_ohlcv(incoming, timeframe)
|
self._validate_ohlcv(incoming)
|
||||||
|
|
||||||
path: Path = self.mk_path(fqme, timeframe)
|
path: Path = self.mk_path(fqme, timeframe)
|
||||||
if path.exists():
|
if path.exists():
|
||||||
stored: pl.DataFrame = self._repair_stored_ohlcv(
|
stored: pl.DataFrame = self._canonicalize_ohlcv(
|
||||||
pl.read_parquet(path),
|
pl.read_parquet(path)
|
||||||
timeframe=timeframe,
|
)
|
||||||
|
stored_len: int = stored.height
|
||||||
|
stored = stored.filter(pl.col('time') > 0)
|
||||||
|
dropped: int = stored_len - stored.height
|
||||||
|
if dropped:
|
||||||
|
log.warning(
|
||||||
|
f'Dropping {dropped} persisted OHLCV row(s) with '
|
||||||
|
f'non-positive timestamps during merge repair'
|
||||||
)
|
)
|
||||||
merged: pl.DataFrame = pl.concat(
|
merged: pl.DataFrame = pl.concat(
|
||||||
[stored, incoming],
|
[stored, incoming],
|
||||||
|
|
@ -588,7 +537,7 @@ class NativeStorageClient:
|
||||||
else:
|
else:
|
||||||
merged = incoming
|
merged = incoming
|
||||||
|
|
||||||
self._validate_ohlcv(merged, timeframe)
|
self._validate_ohlcv(merged)
|
||||||
return self._write_ohlcv(
|
return self._write_ohlcv(
|
||||||
fqme,
|
fqme,
|
||||||
merged,
|
merged,
|
||||||
|
|
|
||||||
|
|
@ -206,34 +206,6 @@ async def notify_backfill(
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
def _synthetic_gap_times(
|
|
||||||
start_t: float,
|
|
||||||
right_t: float|None,
|
|
||||||
count: int,
|
|
||||||
timeframe: float,
|
|
||||||
|
|
||||||
) -> np.ndarray|None:
|
|
||||||
'''
|
|
||||||
Build bounded synthetic timestamps without crossing valid data.
|
|
||||||
|
|
||||||
'''
|
|
||||||
times: np.ndarray = (
|
|
||||||
start_t
|
|
||||||
+
|
|
||||||
np.arange(1, count + 1)*timeframe
|
|
||||||
)
|
|
||||||
if right_t is None:
|
|
||||||
return None
|
|
||||||
|
|
||||||
if not np.isclose(
|
|
||||||
right_t,
|
|
||||||
start_t + (count + 1)*timeframe,
|
|
||||||
):
|
|
||||||
return None
|
|
||||||
|
|
||||||
return times
|
|
||||||
|
|
||||||
|
|
||||||
async def shm_push_in_between(
|
async def shm_push_in_between(
|
||||||
shm: ShmArray,
|
shm: ShmArray,
|
||||||
to_push: np.ndarray,
|
to_push: np.ndarray,
|
||||||
|
|
@ -290,11 +262,11 @@ async def maybe_fill_null_segments(
|
||||||
|
|
||||||
task_status: TaskStatus[None] = trio.TASK_STATUS_IGNORED,
|
task_status: TaskStatus[None] = trio.TASK_STATUS_IGNORED,
|
||||||
|
|
||||||
) -> None:
|
) -> list[Frame]:
|
||||||
|
|
||||||
task_status.started()
|
task_status.started()
|
||||||
|
|
||||||
frame: Frame = shm.array.copy()
|
frame: Frame = shm.array
|
||||||
|
|
||||||
# TODO, put in parent task/daemon root!
|
# TODO, put in parent task/daemon root!
|
||||||
import greenback
|
import greenback
|
||||||
|
|
@ -325,7 +297,6 @@ async def maybe_fill_null_segments(
|
||||||
await tractor.pause()
|
await tractor.pause()
|
||||||
break
|
break
|
||||||
|
|
||||||
try:
|
|
||||||
(
|
(
|
||||||
array,
|
array,
|
||||||
next_start_dt,
|
next_start_dt,
|
||||||
|
|
@ -335,16 +306,6 @@ async def maybe_fill_null_segments(
|
||||||
start_dt=start_dt,
|
start_dt=start_dt,
|
||||||
end_dt=end_dt,
|
end_dt=end_dt,
|
||||||
)
|
)
|
||||||
except (
|
|
||||||
NoData,
|
|
||||||
DataUnavailable,
|
|
||||||
) as exc:
|
|
||||||
log.warning(
|
|
||||||
f'Provider could not repair SHM gap:\n'
|
|
||||||
f'{start_dt} -> {end_dt}\n'
|
|
||||||
f'{exc}\n'
|
|
||||||
)
|
|
||||||
break
|
|
||||||
|
|
||||||
if array.size == 0:
|
if array.size == 0:
|
||||||
log.warning(
|
log.warning(
|
||||||
|
|
@ -353,6 +314,7 @@ async def maybe_fill_null_segments(
|
||||||
)
|
)
|
||||||
# ?TODO? do we want to remove the nulls and push
|
# ?TODO? do we want to remove the nulls and push
|
||||||
# the close price here for the gap duration?
|
# the close price here for the gap duration?
|
||||||
|
await tractor.pause()
|
||||||
break
|
break
|
||||||
|
|
||||||
if (
|
if (
|
||||||
|
|
@ -365,7 +327,7 @@ async def maybe_fill_null_segments(
|
||||||
f'frame_start_dt: {frame_start_dt!r}\n'
|
f'frame_start_dt: {frame_start_dt!r}\n'
|
||||||
f'backfill_until_dt: {backfill_until_dt!r}\n'
|
f'backfill_until_dt: {backfill_until_dt!r}\n'
|
||||||
)
|
)
|
||||||
break
|
await tractor.pause()
|
||||||
|
|
||||||
# XXX TODO: pretty sure if i plot tsla, btcusdt.binance
|
# XXX TODO: pretty sure if i plot tsla, btcusdt.binance
|
||||||
# and mnq.cme.ib this causes a Qt crash XXDDD
|
# and mnq.cme.ib this causes a Qt crash XXDDD
|
||||||
|
|
@ -387,11 +349,21 @@ async def maybe_fill_null_segments(
|
||||||
# - remember that in the display side, only refersh this
|
# - remember that in the display side, only refersh this
|
||||||
# if the respective history is actually "in view".
|
# if the respective history is actually "in view".
|
||||||
# loop
|
# loop
|
||||||
await notify_backfill(
|
try:
|
||||||
sampler_stream,
|
await sampler_stream.send({
|
||||||
mkt,
|
'broadcast_all': {
|
||||||
timeframe,
|
|
||||||
)
|
# XXX NOTE XXX: see the
|
||||||
|
# `.ui._display.increment_history_view()` if block
|
||||||
|
# that looks for this info to FORCE a hard viz
|
||||||
|
# redraw!
|
||||||
|
'backfilling': (mkt.fqme, timeframe),
|
||||||
|
},
|
||||||
|
})
|
||||||
|
except tractor.ContextCancelled:
|
||||||
|
# log.exception
|
||||||
|
await tractor.pause()
|
||||||
|
raise
|
||||||
|
|
||||||
# RECHECK for more null-gaps
|
# RECHECK for more null-gaps
|
||||||
frame: Frame = shm.array
|
frame: Frame = shm.array
|
||||||
|
|
@ -407,7 +379,7 @@ async def maybe_fill_null_segments(
|
||||||
(
|
(
|
||||||
iabs_slices,
|
iabs_slices,
|
||||||
iabs_zero_rows,
|
iabs_zero_rows,
|
||||||
_zero_t,
|
zero_t,
|
||||||
) = null_segs
|
) = null_segs
|
||||||
log.warning(
|
log.warning(
|
||||||
f'{len(iabs_slices)} NULL TIME SEGMENTS DETECTED!\n'
|
f'{len(iabs_slices)} NULL TIME SEGMENTS DETECTED!\n'
|
||||||
|
|
@ -426,61 +398,13 @@ async def maybe_fill_null_segments(
|
||||||
'close',
|
'close',
|
||||||
]
|
]
|
||||||
|
|
||||||
zero_indexes: np.ndarray = np.asarray(iabs_zero_rows)
|
for istart, istop in iabs_slices:
|
||||||
split_at: np.ndarray = (
|
|
||||||
np.flatnonzero(np.diff(zero_indexes) != 1)
|
|
||||||
+
|
|
||||||
1
|
|
||||||
)
|
|
||||||
zero_groups: list[np.ndarray] = np.split(
|
|
||||||
zero_indexes,
|
|
||||||
split_at,
|
|
||||||
)
|
|
||||||
repaired: int = 0
|
|
||||||
frame_start: int = int(frame['index'][0])
|
|
||||||
for zero_group in zero_groups:
|
|
||||||
istart: int = int(zero_group[0])
|
|
||||||
istop: int = int(zero_group[-1]) + 1
|
|
||||||
seed_index: int = istart - 1
|
|
||||||
if seed_index < frame_start:
|
|
||||||
log.warning(
|
|
||||||
f'Can not forward-fill leading SHM nulls:\n'
|
|
||||||
f'{istart}:{istop}\n'
|
|
||||||
)
|
|
||||||
continue
|
|
||||||
|
|
||||||
seed: np.void = shm._array[seed_index]
|
# get view into buffer for null-segment
|
||||||
start_t: float = seed['time']
|
|
||||||
if start_t <= 0:
|
|
||||||
log.warning(
|
|
||||||
f'Can not forward-fill SHM nulls from '
|
|
||||||
f'zero-time seed at {seed_index}\n'
|
|
||||||
)
|
|
||||||
continue
|
|
||||||
|
|
||||||
# Fill only the null rows, preserving both valid margins.
|
|
||||||
gap: np.ndarray = shm._array[istart:istop]
|
gap: np.ndarray = shm._array[istart:istop]
|
||||||
cls: float = seed['close']
|
|
||||||
|
|
||||||
frame_end: int = int(frame['index'][-1])
|
# copy the oldest OHLC samples forward
|
||||||
right_t: float|None = (
|
cls: float = shm._array[istart]['close']
|
||||||
shm._array[istop]['time']
|
|
||||||
if istop <= frame_end
|
|
||||||
else None
|
|
||||||
)
|
|
||||||
gap_times: np.ndarray|None = _synthetic_gap_times(
|
|
||||||
start_t,
|
|
||||||
right_t,
|
|
||||||
len(gap),
|
|
||||||
timeframe,
|
|
||||||
)
|
|
||||||
if gap_times is None:
|
|
||||||
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?
|
# TODO: how can we mark this range as being a gap tho?
|
||||||
# -[ ] maybe pg finally supports nulls in ndarray to
|
# -[ ] maybe pg finally supports nulls in ndarray to
|
||||||
|
|
@ -489,13 +413,30 @@ async def maybe_fill_null_segments(
|
||||||
# another col/field to denote?
|
# another col/field to denote?
|
||||||
gap[ohlc_fields] = cls
|
gap[ohlc_fields] = cls
|
||||||
|
|
||||||
gap['time'] = gap_times
|
start_t: float = shm._array[istart]['time']
|
||||||
repaired += len(gap)
|
t_diff: float = (istop - istart)*timeframe
|
||||||
|
|
||||||
|
gap['time'] = np.arange(
|
||||||
|
start=start_t,
|
||||||
|
stop=start_t + t_diff,
|
||||||
|
step=timeframe,
|
||||||
|
)
|
||||||
|
|
||||||
# TODO: reimpl using the new `.ui._remote_ctl` ctx
|
# TODO: reimpl using the new `.ui._remote_ctl` ctx
|
||||||
# ideally using some kinda decent
|
# ideally using some kinda decent
|
||||||
# tractory-reverse-lookup-connnection from some other
|
# tractory-reverse-lookup-connnection from some other
|
||||||
# `Context` type thingy?
|
# `Context` type thingy?
|
||||||
|
await sampler_stream.send({
|
||||||
|
'broadcast_all': {
|
||||||
|
|
||||||
|
# XXX NOTE XXX: see the
|
||||||
|
# `.ui._display.increment_history_view()` if block
|
||||||
|
# that looks for this info to FORCE a hard viz
|
||||||
|
# redraw!
|
||||||
|
'backfilling': (mkt.fqme, timeframe),
|
||||||
|
},
|
||||||
|
})
|
||||||
|
|
||||||
# TODO: interatively step through any remaining
|
# TODO: interatively step through any remaining
|
||||||
# time-gaps/null-segments and spawn piecewise backfiller
|
# time-gaps/null-segments and spawn piecewise backfiller
|
||||||
# tasks in a nursery?
|
# tasks in a nursery?
|
||||||
|
|
@ -506,13 +447,6 @@ async def maybe_fill_null_segments(
|
||||||
# -[ ] fill algo: do queries in alternating "latest, then
|
# -[ ] fill algo: do queries in alternating "latest, then
|
||||||
# earliest, then latest.. etc?"
|
# earliest, then latest.. etc?"
|
||||||
|
|
||||||
if repaired:
|
|
||||||
await notify_backfill(
|
|
||||||
sampler_stream,
|
|
||||||
mkt,
|
|
||||||
timeframe,
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
async def start_backfill(
|
async def start_backfill(
|
||||||
get_hist,
|
get_hist,
|
||||||
|
|
@ -746,12 +680,7 @@ async def start_backfill(
|
||||||
array,
|
array,
|
||||||
prepend_until_dt=backfill_until_dt,
|
prepend_until_dt=backfill_until_dt,
|
||||||
)
|
)
|
||||||
# The requested end bar is already the oldest published
|
to_push: np.ndarray = to_store
|
||||||
# sample. Keep it for idempotent storage merge, but do not
|
|
||||||
# consume another physical SHM row for the overlap.
|
|
||||||
to_push: np.ndarray = to_store[
|
|
||||||
to_store['time'] < last_start_dt.timestamp()
|
|
||||||
]
|
|
||||||
ln: int = len(to_push)
|
ln: int = len(to_push)
|
||||||
if ln:
|
if ln:
|
||||||
log.info(
|
log.info(
|
||||||
|
|
@ -1385,25 +1314,27 @@ async def tsdb_backfill(
|
||||||
|
|
||||||
log.info(f'Loaded {to_push.shape} datums from storage')
|
log.info(f'Loaded {to_push.shape} datums from storage')
|
||||||
|
|
||||||
# 2nd nursery END
|
# NOTE: ASYNC-conduct tsdb timestamp gap detection and backfill any
|
||||||
|
# seemingly missing (null-time) segments..
|
||||||
if last_tsdb_dt is not None:
|
# TODO: ideally these can never exist!
|
||||||
# Repair only after reverse backfill releases providers
|
|
||||||
# such as IB which permit one history request at a time.
|
|
||||||
# TODO: ideally these null segments can never exist!
|
|
||||||
# -[ ] somehow it seems sometimes we're writing zero-ed
|
# -[ ] somehow it seems sometimes we're writing zero-ed
|
||||||
# segments to tsdbs during teardown?
|
# segments to tsdbs during teardown?
|
||||||
# -[ ] can we ensure that the backfiller tasks do this
|
# -[ ] can we ensure that the backfiller tasks do this
|
||||||
# work PREVENTAVELY instead?
|
# work PREVENTAVELY instead?
|
||||||
# -[ ] fill in non-zero epoch time values ALWAYS!
|
# -[ ] fill in non-zero epoch time values ALWAYS!
|
||||||
await maybe_fill_null_segments(
|
# await maybe_fill_null_segments(
|
||||||
|
await tn.start(partial(
|
||||||
|
maybe_fill_null_segments,
|
||||||
|
|
||||||
shm=shm,
|
shm=shm,
|
||||||
timeframe=timeframe,
|
timeframe=timeframe,
|
||||||
get_hist=get_hist,
|
get_hist=get_hist,
|
||||||
sampler_stream=sampler_stream,
|
sampler_stream=sampler_stream,
|
||||||
mkt=mkt,
|
mkt=mkt,
|
||||||
backfill_until_dt=last_tsdb_dt,
|
backfill_until_dt=last_tsdb_dt,
|
||||||
)
|
))
|
||||||
|
|
||||||
|
# 2nd nursery END
|
||||||
|
|
||||||
# TODO: maybe start history anal and load missing "history
|
# TODO: maybe start history anal and load missing "history
|
||||||
# gaps" via backend..
|
# gaps" via backend..
|
||||||
|
|
|
||||||
|
|
@ -29,8 +29,6 @@ from piker.storage.nativedb import (
|
||||||
ohlc_key_map,
|
ohlc_key_map,
|
||||||
)
|
)
|
||||||
from piker.tsp._history import (
|
from piker.tsp._history import (
|
||||||
_synthetic_gap_times,
|
|
||||||
maybe_fill_null_segments,
|
|
||||||
notify_backfill,
|
notify_backfill,
|
||||||
publish_latest_frame,
|
publish_latest_frame,
|
||||||
start_backfill,
|
start_backfill,
|
||||||
|
|
@ -94,22 +92,15 @@ def test_empty_frame_completes() -> None:
|
||||||
|
|
||||||
def test_storage_receives_full_provider_delta() -> None:
|
def test_storage_receives_full_provider_delta() -> None:
|
||||||
'''
|
'''
|
||||||
SHM drops the published endpoint without truncating storage.
|
SHM capacity truncation does not truncate durable history.
|
||||||
|
|
||||||
Reverse history queries include their requested end bar, which is
|
|
||||||
already the oldest sample in SHM. Publishing the overlap consumed
|
|
||||||
an extra physical row and made reserved null counts disagree with
|
|
||||||
elapsed timestamps. Return a frame with the overlap and prove
|
|
||||||
NativeDB receives the complete provider delta while SHM receives
|
|
||||||
only bars older than its published boundary.
|
|
||||||
|
|
||||||
'''
|
'''
|
||||||
frame = np.zeros(
|
frame = np.zeros(
|
||||||
3,
|
2,
|
||||||
dtype=np.dtype(def_iohlcv_fields),
|
dtype=np.dtype(def_iohlcv_fields),
|
||||||
)
|
)
|
||||||
frame['index'] = [0, 1, 2]
|
frame['index'] = [0, 1]
|
||||||
frame['time'] = [60, 120, 180]
|
frame['time'] = [60, 120]
|
||||||
|
|
||||||
events: list[str] = []
|
events: list[str] = []
|
||||||
|
|
||||||
|
|
@ -151,7 +142,7 @@ def test_storage_receives_full_provider_delta() -> None:
|
||||||
return (
|
return (
|
||||||
frame,
|
frame,
|
||||||
from_timestamp(60),
|
from_timestamp(60),
|
||||||
from_timestamp(180),
|
from_timestamp(120),
|
||||||
)
|
)
|
||||||
|
|
||||||
async def main() -> None:
|
async def main() -> None:
|
||||||
|
|
@ -163,7 +154,7 @@ def test_storage_receives_full_provider_delta() -> None:
|
||||||
mkt=mkt,
|
mkt=mkt,
|
||||||
shm=shm,
|
shm=shm,
|
||||||
timeframe=60,
|
timeframe=60,
|
||||||
backfill_from_shm_index=2,
|
backfill_from_shm_index=1,
|
||||||
backfill_from_dt=from_timestamp(180),
|
backfill_from_dt=from_timestamp(180),
|
||||||
sampler_stream=Sampler(),
|
sampler_stream=Sampler(),
|
||||||
backfill_until_dt=from_timestamp(60),
|
backfill_until_dt=from_timestamp(60),
|
||||||
|
|
@ -172,8 +163,8 @@ def test_storage_receives_full_provider_delta() -> None:
|
||||||
)
|
)
|
||||||
|
|
||||||
trio.run(main)
|
trio.run(main)
|
||||||
assert shm.pushed[0]['time'].tolist() == [60, 120]
|
assert shm.pushed[0]['time'].tolist() == [120]
|
||||||
assert storage.frames[0]['time'].tolist() == [60, 120, 180]
|
assert storage.frames[0]['time'].tolist() == [60, 120]
|
||||||
assert events == ['storage', 'shm', 'sampler']
|
assert events == ['storage', 'shm', 'sampler']
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -353,163 +344,3 @@ def test_backfill_notification_timeout_is_bounded(
|
||||||
)
|
)
|
||||||
|
|
||||||
trio.run(main)
|
trio.run(main)
|
||||||
|
|
||||||
|
|
||||||
def test_synthetic_gap_times_require_valid_right_boundary() -> None:
|
|
||||||
'''
|
|
||||||
Synthetic rows must remain strictly between valid boundary bars.
|
|
||||||
|
|
||||||
The live QQQ reservation held one more physical zero row than its
|
|
||||||
timestamp interval could represent. Blind fill would assign the
|
|
||||||
null the right boundary's timestamp, creating a duplicate. Prove
|
|
||||||
aligned bounds produce cadence while an overfull segment
|
|
||||||
is rejected instead of inventing non-monotonic timestamps.
|
|
||||||
|
|
||||||
'''
|
|
||||||
aligned = _synthetic_gap_times(
|
|
||||||
start_t=60,
|
|
||||||
right_t=300,
|
|
||||||
count=3,
|
|
||||||
timeframe=60,
|
|
||||||
)
|
|
||||||
overfull = _synthetic_gap_times(
|
|
||||||
start_t=60,
|
|
||||||
right_t=240,
|
|
||||||
count=3,
|
|
||||||
timeframe=60,
|
|
||||||
)
|
|
||||||
trailing = _synthetic_gap_times(
|
|
||||||
start_t=60,
|
|
||||||
right_t=None,
|
|
||||||
count=3,
|
|
||||||
timeframe=60,
|
|
||||||
)
|
|
||||||
|
|
||||||
assert aligned is not None
|
|
||||||
assert aligned.tolist() == [120, 180, 240]
|
|
||||||
assert overfull is None
|
|
||||||
assert trailing is None
|
|
||||||
|
|
||||||
|
|
||||||
def test_null_repair_falls_back_without_debug_pause(
|
|
||||||
monkeypatch: pytest.MonkeyPatch,
|
|
||||||
) -> None:
|
|
||||||
'''
|
|
||||||
Empty provider repair must resolve interior zero rows locally.
|
|
||||||
|
|
||||||
QQQ startup exposed 653 zero-time rows between stored history and
|
|
||||||
an IB frame. The null repair queried IB while reverse backfill
|
|
||||||
occupied its sole history request slot, then could
|
|
||||||
enter an interactive debug pause before reaching forward-fill.
|
|
||||||
Arrange an interior null run, return an empty provider frame, 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.
|
|
||||||
|
|
||||||
'''
|
|
||||||
monkeypatch.setattr(
|
|
||||||
'piker.tsp._history._notify_timeout_s',
|
|
||||||
0.01,
|
|
||||||
)
|
|
||||||
|
|
||||||
backing = np.zeros(
|
|
||||||
105,
|
|
||||||
dtype=np.dtype(def_iohlcv_fields),
|
|
||||||
)
|
|
||||||
backing['index'] = np.arange(105)
|
|
||||||
frame = backing[100:105]
|
|
||||||
frame['index'] = np.arange(100, 105)
|
|
||||||
frame['time'] = [60, 0, 0, 0, 300]
|
|
||||||
frame['open'] = [9, 0, 0, 0, 19]
|
|
||||||
frame['high'] = [11, 0, 0, 0, 21]
|
|
||||||
frame['low'] = [8, 0, 0, 0, 18]
|
|
||||||
frame['close'] = [10, 0, 0, 0, 20]
|
|
||||||
frame['volume'] = [5, 0, 0, 0, 6]
|
|
||||||
|
|
||||||
class Shm:
|
|
||||||
'''
|
|
||||||
Expose the test frame through the SHM repair interface.
|
|
||||||
|
|
||||||
'''
|
|
||||||
def __init__(self) -> None:
|
|
||||||
self._array = backing
|
|
||||||
|
|
||||||
@property
|
|
||||||
def array(self) -> np.ndarray:
|
|
||||||
'''
|
|
||||||
Return the readable SHM view.
|
|
||||||
|
|
||||||
'''
|
|
||||||
return self._array[100:105]
|
|
||||||
|
|
||||||
class Sampler:
|
|
||||||
'''
|
|
||||||
Model a backpressured UI notification stream.
|
|
||||||
|
|
||||||
'''
|
|
||||||
async def send(self, msg) -> None:
|
|
||||||
'''
|
|
||||||
Block until the bounded notifier cancels this send.
|
|
||||||
|
|
||||||
'''
|
|
||||||
await trio.sleep_forever()
|
|
||||||
|
|
||||||
async def get_hist(*args, **kwargs):
|
|
||||||
'''
|
|
||||||
Return no provider bars for the requested null segment.
|
|
||||||
|
|
||||||
'''
|
|
||||||
end_dt = kwargs['end_dt']
|
|
||||||
empty = np.empty(
|
|
||||||
0,
|
|
||||||
dtype=np.dtype(def_iohlcv_fields),
|
|
||||||
)
|
|
||||||
return empty, end_dt, end_dt
|
|
||||||
|
|
||||||
async def fail_pause() -> None:
|
|
||||||
'''
|
|
||||||
Reject interactive debugging in normal repair fallbacks.
|
|
||||||
|
|
||||||
'''
|
|
||||||
raise AssertionError('null repair entered tractor.pause()')
|
|
||||||
|
|
||||||
monkeypatch.setattr(tractor, 'pause', fail_pause)
|
|
||||||
|
|
||||||
async def main() -> None:
|
|
||||||
import greenback
|
|
||||||
|
|
||||||
async def ensure_portal() -> None:
|
|
||||||
'''
|
|
||||||
Avoid installing a greenback portal in this unit test.
|
|
||||||
|
|
||||||
'''
|
|
||||||
|
|
||||||
monkeypatch.setattr(
|
|
||||||
greenback,
|
|
||||||
'ensure_portal',
|
|
||||||
ensure_portal,
|
|
||||||
)
|
|
||||||
with trio.fail_after(0.1):
|
|
||||||
await maybe_fill_null_segments(
|
|
||||||
shm=Shm(),
|
|
||||||
timeframe=60,
|
|
||||||
get_hist=get_hist,
|
|
||||||
sampler_stream=Sampler(),
|
|
||||||
mkt=SimpleNamespace(fqme='qqq.nasdaq.ib'),
|
|
||||||
backfill_until_dt=from_timestamp(60),
|
|
||||||
)
|
|
||||||
|
|
||||||
trio.run(main)
|
|
||||||
|
|
||||||
assert frame['time'].tolist() == [60, 120, 180, 240, 300]
|
|
||||||
for field in ('open', 'high', 'low', 'close'):
|
|
||||||
assert frame[field].tolist() == [
|
|
||||||
frame[field][0],
|
|
||||||
10,
|
|
||||||
10,
|
|
||||||
10,
|
|
||||||
frame[field][-1],
|
|
||||||
]
|
|
||||||
assert frame['volume'].tolist() == [5, 0, 0, 0, 6]
|
|
||||||
|
|
|
||||||
|
|
@ -8,19 +8,11 @@ from datetime import (
|
||||||
)
|
)
|
||||||
from types import SimpleNamespace
|
from types import SimpleNamespace
|
||||||
|
|
||||||
from ib_async import RequestError
|
from pendulum import datetime as pdatetime
|
||||||
import numpy as np
|
|
||||||
from pendulum import (
|
|
||||||
datetime as pdatetime,
|
|
||||||
duration,
|
|
||||||
)
|
|
||||||
import pytest
|
import pytest
|
||||||
import trio
|
import trio
|
||||||
|
|
||||||
from piker.brokers import DataUnavailable
|
|
||||||
from piker.brokers.ib import feed as ib_feed
|
|
||||||
from piker.brokers.ib.api import Client
|
from piker.brokers.ib.api import Client
|
||||||
from piker.brokers.ib.feed import get_bars
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.parametrize(
|
@pytest.mark.parametrize(
|
||||||
|
|
@ -87,502 +79,3 @@ def test_history_end_datetime_uses_ib_utc_format(
|
||||||
|
|
||||||
trio.run(main)
|
trio.run(main)
|
||||||
assert request['endDateTime'] == expected
|
assert request['endDateTime'] == expected
|
||||||
|
|
||||||
|
|
||||||
def test_empty_history_frame_completes() -> None:
|
|
||||||
'''
|
|
||||||
A blank IB response must release the history request slot.
|
|
||||||
|
|
||||||
``get_bars()`` previously returned from its query task without
|
|
||||||
signalling the parent, which then reset the feed forever. That
|
|
||||||
kept serialized SHM repair from running. Return IB's empty
|
|
||||||
response and prove the request completes explicitly
|
|
||||||
with no result instead of relying on a reset timeout.
|
|
||||||
|
|
||||||
'''
|
|
||||||
class Proxy:
|
|
||||||
'''
|
|
||||||
Return one empty IB history response.
|
|
||||||
|
|
||||||
'''
|
|
||||||
async def bars(self, **kwargs):
|
|
||||||
'''
|
|
||||||
Model ``Client.bars()`` during a venue data gap.
|
|
||||||
|
|
||||||
'''
|
|
||||||
return (
|
|
||||||
[],
|
|
||||||
np.empty(0),
|
|
||||||
duration(seconds=2000),
|
|
||||||
)
|
|
||||||
|
|
||||||
async def main() -> None:
|
|
||||||
with trio.fail_after(0.1):
|
|
||||||
result, timedout = await get_bars(
|
|
||||||
Proxy(),
|
|
||||||
'qqq.nasdaq.ib',
|
|
||||||
1,
|
|
||||||
)
|
|
||||||
assert result is None
|
|
||||||
assert not timedout
|
|
||||||
|
|
||||||
trio.run(main)
|
|
||||||
|
|
||||||
|
|
||||||
def test_repeated_ib_cancellation_is_bounded() -> None:
|
|
||||||
'''
|
|
||||||
Repeated IB cancellation must not monopolize history forever.
|
|
||||||
|
|
||||||
IB can answer a historical request with error 162 cancellation.
|
|
||||||
The old retry loop had no limit, so reverse backfill could retain
|
|
||||||
the backend's sole request slot and prevent final null repair.
|
|
||||||
Always raise that response and prove the configured retry bound
|
|
||||||
terminates with ``DataUnavailable`` deterministically.
|
|
||||||
|
|
||||||
'''
|
|
||||||
class Proxy:
|
|
||||||
'''
|
|
||||||
Cancel every IB history request.
|
|
||||||
|
|
||||||
'''
|
|
||||||
async def bars(self, **kwargs):
|
|
||||||
'''
|
|
||||||
Raise IB's historical-query cancellation response.
|
|
||||||
|
|
||||||
'''
|
|
||||||
raise RequestError(
|
|
||||||
1,
|
|
||||||
162,
|
|
||||||
'API historical data query cancelled',
|
|
||||||
)
|
|
||||||
|
|
||||||
async def main() -> None:
|
|
||||||
with (
|
|
||||||
trio.fail_after(0.1),
|
|
||||||
pytest.raises(DataUnavailable),
|
|
||||||
):
|
|
||||||
await get_bars(
|
|
||||||
Proxy(),
|
|
||||||
'qqq.nasdaq.ib',
|
|
||||||
1,
|
|
||||||
max_failed_resets=2,
|
|
||||||
)
|
|
||||||
|
|
||||||
trio.run(main)
|
|
||||||
|
|
||||||
|
|
||||||
def test_query_error_cancels_stalled_reset(
|
|
||||||
monkeypatch: pytest.MonkeyPatch,
|
|
||||||
) -> None:
|
|
||||||
'''
|
|
||||||
Exceptional history exit must cancel its in-flight reset task.
|
|
||||||
|
|
||||||
A timed-out request acquires the reset lock before launching
|
|
||||||
IB's VNC reset task. If the query then raises while that task is
|
|
||||||
blocked, stale work must not retain the reset lock. Force that
|
|
||||||
ordering with explicit Trio checkpoints, then prove bounded query
|
|
||||||
failure cancels the reset task before propagating.
|
|
||||||
|
|
||||||
'''
|
|
||||||
reset_cancelled: list[bool] = []
|
|
||||||
|
|
||||||
class Proxy:
|
|
||||||
'''
|
|
||||||
Delay once, then return IB cancellation.
|
|
||||||
|
|
||||||
'''
|
|
||||||
async def bars(self, **kwargs):
|
|
||||||
'''
|
|
||||||
Let the parent acquire the reset lock before failing.
|
|
||||||
|
|
||||||
'''
|
|
||||||
await trio.sleep(0.02)
|
|
||||||
raise RequestError(
|
|
||||||
1,
|
|
||||||
162,
|
|
||||||
'API historical data query cancelled',
|
|
||||||
)
|
|
||||||
|
|
||||||
async def stalled_reset(
|
|
||||||
proxy,
|
|
||||||
reset_type='data',
|
|
||||||
task_status=trio.TASK_STATUS_IGNORED,
|
|
||||||
) -> bool:
|
|
||||||
'''
|
|
||||||
Publish reset state, then block until sibling failure.
|
|
||||||
|
|
||||||
'''
|
|
||||||
done = trio.Event()
|
|
||||||
result: list[bool] = []
|
|
||||||
with trio.CancelScope() as cs:
|
|
||||||
task_status.started((cs, done, result))
|
|
||||||
try:
|
|
||||||
await trio.sleep_forever()
|
|
||||||
finally:
|
|
||||||
reset_cancelled.append(True)
|
|
||||||
return False
|
|
||||||
|
|
||||||
monkeypatch.setattr(
|
|
||||||
ib_feed,
|
|
||||||
'wait_on_data_reset',
|
|
||||||
stalled_reset,
|
|
||||||
)
|
|
||||||
|
|
||||||
async def main() -> None:
|
|
||||||
with (
|
|
||||||
trio.fail_after(0.1),
|
|
||||||
pytest.raises(DataUnavailable),
|
|
||||||
):
|
|
||||||
await get_bars(
|
|
||||||
Proxy(),
|
|
||||||
'qqq.nasdaq.ib',
|
|
||||||
1,
|
|
||||||
feed_reset_timeout=0.001,
|
|
||||||
max_failed_resets=1,
|
|
||||||
)
|
|
||||||
|
|
||||||
trio.run(main)
|
|
||||||
assert reset_cancelled == [True]
|
|
||||||
|
|
||||||
|
|
||||||
def test_failed_pacing_resets_reach_retry_limit(
|
|
||||||
monkeypatch: pytest.MonkeyPatch,
|
|
||||||
) -> None:
|
|
||||||
'''
|
|
||||||
Failed pacing resets must increment the request's retry count.
|
|
||||||
|
|
||||||
The reset API returns an ``Event`` and result box. The old
|
|
||||||
code tested the always-truthy event and reset the failure counter
|
|
||||||
after each failure. Return immediate failures and prove two
|
|
||||||
pacing responses terminate instead of looping forever.
|
|
||||||
|
|
||||||
'''
|
|
||||||
reset_calls: list[bool] = []
|
|
||||||
|
|
||||||
class Proxy:
|
|
||||||
'''
|
|
||||||
Return a pacing violation for every request.
|
|
||||||
|
|
||||||
'''
|
|
||||||
_aio_ns = SimpleNamespace(
|
|
||||||
ib=SimpleNamespace(
|
|
||||||
client=SimpleNamespace(),
|
|
||||||
),
|
|
||||||
)
|
|
||||||
|
|
||||||
async def bars(self, **kwargs):
|
|
||||||
'''
|
|
||||||
Raise IB's pacing response.
|
|
||||||
|
|
||||||
'''
|
|
||||||
raise RequestError(
|
|
||||||
1,
|
|
||||||
162,
|
|
||||||
ib_feed._pacing,
|
|
||||||
)
|
|
||||||
|
|
||||||
async def failed_reset(
|
|
||||||
proxy,
|
|
||||||
reset_type='data',
|
|
||||||
task_status=trio.TASK_STATUS_IGNORED,
|
|
||||||
) -> bool:
|
|
||||||
'''
|
|
||||||
Publish one completed, unsuccessful reset.
|
|
||||||
|
|
||||||
'''
|
|
||||||
done = trio.Event()
|
|
||||||
result = [False]
|
|
||||||
reset_calls.append(True)
|
|
||||||
with trio.CancelScope() as cs:
|
|
||||||
task_status.started((cs, done, result))
|
|
||||||
done.set()
|
|
||||||
return False
|
|
||||||
|
|
||||||
monkeypatch.setattr(
|
|
||||||
ib_feed,
|
|
||||||
'wait_on_data_reset',
|
|
||||||
failed_reset,
|
|
||||||
)
|
|
||||||
|
|
||||||
async def main() -> None:
|
|
||||||
with (
|
|
||||||
trio.fail_after(0.1),
|
|
||||||
pytest.raises(DataUnavailable),
|
|
||||||
):
|
|
||||||
await get_bars(
|
|
||||||
Proxy(),
|
|
||||||
'qqq.nasdaq.ib',
|
|
||||||
1,
|
|
||||||
max_failed_resets=2,
|
|
||||||
)
|
|
||||||
|
|
||||||
trio.run(main)
|
|
||||||
assert reset_calls == [True, True]
|
|
||||||
|
|
||||||
|
|
||||||
def test_pacing_resets_are_serialized_across_requests(
|
|
||||||
monkeypatch: pytest.MonkeyPatch,
|
|
||||||
) -> None:
|
|
||||||
'''
|
|
||||||
Concurrent history requests must never reset IB farms together.
|
|
||||||
|
|
||||||
Pacing recovery bypassed reset ownership, allowing 1s and 60s
|
|
||||||
requests to launch the VNC hack concurrently. Start two pacing
|
|
||||||
failures in one nursery and hold each fake reset across a Trio
|
|
||||||
checkpoint. Active-count assertions prove the actor-global lock
|
|
||||||
serializes reset lifetimes without depending on scheduler timing.
|
|
||||||
|
|
||||||
'''
|
|
||||||
active: int = 0
|
|
||||||
max_active: int = 0
|
|
||||||
reset_calls: int = 0
|
|
||||||
|
|
||||||
class Proxy:
|
|
||||||
'''
|
|
||||||
Return a pacing violation for every request.
|
|
||||||
|
|
||||||
'''
|
|
||||||
_aio_ns = SimpleNamespace(
|
|
||||||
ib=SimpleNamespace(
|
|
||||||
client=SimpleNamespace(),
|
|
||||||
),
|
|
||||||
)
|
|
||||||
|
|
||||||
async def bars(self, **kwargs):
|
|
||||||
'''
|
|
||||||
Raise IB's pacing response.
|
|
||||||
|
|
||||||
'''
|
|
||||||
raise RequestError(
|
|
||||||
1,
|
|
||||||
162,
|
|
||||||
ib_feed._pacing,
|
|
||||||
)
|
|
||||||
|
|
||||||
async def failed_reset(
|
|
||||||
proxy,
|
|
||||||
reset_type='data',
|
|
||||||
task_status=trio.TASK_STATUS_IGNORED,
|
|
||||||
) -> bool:
|
|
||||||
'''
|
|
||||||
Hold one failed reset long enough to expose overlap.
|
|
||||||
|
|
||||||
'''
|
|
||||||
nonlocal active, max_active, reset_calls
|
|
||||||
done = trio.Event()
|
|
||||||
result: list[bool] = []
|
|
||||||
with trio.CancelScope() as cs:
|
|
||||||
task_status.started((cs, done, result))
|
|
||||||
active += 1
|
|
||||||
reset_calls += 1
|
|
||||||
max_active = max(max_active, active)
|
|
||||||
await trio.sleep(0.01)
|
|
||||||
active -= 1
|
|
||||||
result.append(False)
|
|
||||||
done.set()
|
|
||||||
return False
|
|
||||||
|
|
||||||
monkeypatch.setattr(
|
|
||||||
ib_feed,
|
|
||||||
'wait_on_data_reset',
|
|
||||||
failed_reset,
|
|
||||||
)
|
|
||||||
|
|
||||||
async def request() -> None:
|
|
||||||
'''
|
|
||||||
Run one independently bounded history request.
|
|
||||||
|
|
||||||
'''
|
|
||||||
with pytest.raises(DataUnavailable):
|
|
||||||
await get_bars(
|
|
||||||
Proxy(),
|
|
||||||
'qqq.nasdaq.ib',
|
|
||||||
1,
|
|
||||||
max_failed_resets=1,
|
|
||||||
)
|
|
||||||
|
|
||||||
async def main() -> None:
|
|
||||||
with trio.fail_after(0.1):
|
|
||||||
async with trio.open_nursery() as nursery:
|
|
||||||
nursery.start_soon(request)
|
|
||||||
nursery.start_soon(request)
|
|
||||||
|
|
||||||
trio.run(main)
|
|
||||||
assert reset_calls == 2
|
|
||||||
assert max_active == 1
|
|
||||||
|
|
||||||
|
|
||||||
def test_pacing_supersedes_timeout_reset_once(
|
|
||||||
monkeypatch: pytest.MonkeyPatch,
|
|
||||||
) -> None:
|
|
||||||
'''
|
|
||||||
Pacing handoff must not count cancellation as another failure.
|
|
||||||
|
|
||||||
A slow request can trigger timeout recovery immediately before IB
|
|
||||||
returns a pacing violation. Query-side recovery then deliberately
|
|
||||||
cancels that data reset and waits for the same actor-global lock.
|
|
||||||
Reproduce this with blocked first and failed second resets. A
|
|
||||||
one-attempt limit proves cancellation releases the lock
|
|
||||||
without aborting before the intended connection reset runs.
|
|
||||||
|
|
||||||
'''
|
|
||||||
reset_types: list[str] = []
|
|
||||||
|
|
||||||
class Proxy:
|
|
||||||
'''
|
|
||||||
Return pacing only after timeout recovery starts.
|
|
||||||
|
|
||||||
'''
|
|
||||||
_aio_ns = SimpleNamespace(
|
|
||||||
ib=SimpleNamespace(
|
|
||||||
client=SimpleNamespace(),
|
|
||||||
),
|
|
||||||
)
|
|
||||||
|
|
||||||
async def bars(self, **kwargs):
|
|
||||||
'''
|
|
||||||
Delay past the reset timeout, then report pacing.
|
|
||||||
|
|
||||||
'''
|
|
||||||
await trio.sleep(0.02)
|
|
||||||
raise RequestError(
|
|
||||||
1,
|
|
||||||
162,
|
|
||||||
ib_feed._pacing,
|
|
||||||
)
|
|
||||||
|
|
||||||
async def reset(
|
|
||||||
proxy,
|
|
||||||
reset_type='data',
|
|
||||||
task_status=trio.TASK_STATUS_IGNORED,
|
|
||||||
) -> bool:
|
|
||||||
'''
|
|
||||||
Block data reset and fail connection reset immediately.
|
|
||||||
|
|
||||||
'''
|
|
||||||
done = trio.Event()
|
|
||||||
result: list[bool] = []
|
|
||||||
reset_types.append(reset_type)
|
|
||||||
with trio.CancelScope() as cs:
|
|
||||||
task_status.started((cs, done, result))
|
|
||||||
if reset_type == 'data':
|
|
||||||
await trio.sleep_forever()
|
|
||||||
|
|
||||||
result.append(False)
|
|
||||||
done.set()
|
|
||||||
return False
|
|
||||||
|
|
||||||
monkeypatch.setattr(
|
|
||||||
ib_feed,
|
|
||||||
'wait_on_data_reset',
|
|
||||||
reset,
|
|
||||||
)
|
|
||||||
|
|
||||||
async def main() -> None:
|
|
||||||
with (
|
|
||||||
trio.fail_after(0.1),
|
|
||||||
pytest.raises(DataUnavailable),
|
|
||||||
):
|
|
||||||
await get_bars(
|
|
||||||
Proxy(),
|
|
||||||
'qqq.nasdaq.ib',
|
|
||||||
1,
|
|
||||||
feed_reset_timeout=0.001,
|
|
||||||
max_failed_resets=1,
|
|
||||||
)
|
|
||||||
|
|
||||||
trio.run(main)
|
|
||||||
assert reset_types == ['data', 'connection']
|
|
||||||
|
|
||||||
|
|
||||||
def test_completed_pacing_reset_restarts_timeout_wait(
|
|
||||||
monkeypatch: pytest.MonkeyPatch,
|
|
||||||
) -> None:
|
|
||||||
'''
|
|
||||||
A queued timeout waiter must consume completed pacing recovery.
|
|
||||||
|
|
||||||
While query-side pacing recovery holds the reset lock, the outer
|
|
||||||
timeout task can already queue behind it. Without a completion
|
|
||||||
handoff, that waiter launches a redundant data reset immediately.
|
|
||||||
Return a blank frame just after a successful connection reset and
|
|
||||||
prove no third reset runs before the restarted timeout observes
|
|
||||||
completion.
|
|
||||||
|
|
||||||
'''
|
|
||||||
reset_types: list[str] = []
|
|
||||||
|
|
||||||
class Proxy:
|
|
||||||
'''
|
|
||||||
Pace the first query and return blank data on retry.
|
|
||||||
|
|
||||||
'''
|
|
||||||
_aio_ns = SimpleNamespace(
|
|
||||||
ib=SimpleNamespace(
|
|
||||||
client=SimpleNamespace(),
|
|
||||||
),
|
|
||||||
)
|
|
||||||
|
|
||||||
def __init__(self) -> None:
|
|
||||||
self.calls: int = 0
|
|
||||||
|
|
||||||
async def bars(self, **kwargs):
|
|
||||||
'''
|
|
||||||
Control the pacing-to-success handoff ordering.
|
|
||||||
|
|
||||||
'''
|
|
||||||
self.calls += 1
|
|
||||||
if self.calls == 1:
|
|
||||||
await trio.sleep(0.02)
|
|
||||||
raise RequestError(
|
|
||||||
1,
|
|
||||||
162,
|
|
||||||
ib_feed._pacing,
|
|
||||||
)
|
|
||||||
|
|
||||||
await trio.sleep(0.0005)
|
|
||||||
return (
|
|
||||||
[],
|
|
||||||
np.empty(0),
|
|
||||||
duration(seconds=2000),
|
|
||||||
)
|
|
||||||
|
|
||||||
async def reset(
|
|
||||||
proxy,
|
|
||||||
reset_type='data',
|
|
||||||
task_status=trio.TASK_STATUS_IGNORED,
|
|
||||||
) -> bool:
|
|
||||||
'''
|
|
||||||
Cancel timeout recovery and complete pacing recovery.
|
|
||||||
|
|
||||||
'''
|
|
||||||
done = trio.Event()
|
|
||||||
result: list[bool] = []
|
|
||||||
reset_types.append(reset_type)
|
|
||||||
with trio.CancelScope() as cs:
|
|
||||||
task_status.started((cs, done, result))
|
|
||||||
if reset_type == 'data':
|
|
||||||
await trio.sleep_forever()
|
|
||||||
|
|
||||||
result.append(reset_type == 'connection')
|
|
||||||
done.set()
|
|
||||||
return result[0]
|
|
||||||
|
|
||||||
monkeypatch.setattr(
|
|
||||||
ib_feed,
|
|
||||||
'wait_on_data_reset',
|
|
||||||
reset,
|
|
||||||
)
|
|
||||||
|
|
||||||
async def main() -> None:
|
|
||||||
with trio.fail_after(0.1):
|
|
||||||
result, _ = await get_bars(
|
|
||||||
Proxy(),
|
|
||||||
'qqq.nasdaq.ib',
|
|
||||||
1,
|
|
||||||
feed_reset_timeout=0.001,
|
|
||||||
max_failed_resets=2,
|
|
||||||
)
|
|
||||||
assert result is None
|
|
||||||
|
|
||||||
trio.run(main)
|
|
||||||
assert reset_types == ['data', 'connection']
|
|
||||||
|
|
|
||||||
|
|
@ -1,235 +0,0 @@
|
||||||
'''
|
|
||||||
Shared-memory persistence regressions.
|
|
||||||
|
|
||||||
'''
|
|
||||||
from collections.abc import (
|
|
||||||
AsyncIterator,
|
|
||||||
Iterator,
|
|
||||||
)
|
|
||||||
from contextlib import asynccontextmanager
|
|
||||||
from pathlib import Path
|
|
||||||
from types import SimpleNamespace
|
|
||||||
from typing import Any
|
|
||||||
|
|
||||||
import numpy as np
|
|
||||||
import polars as pl
|
|
||||||
import pytest
|
|
||||||
|
|
||||||
from piker import tsp
|
|
||||||
from piker.storage import cli as storage_cli
|
|
||||||
from piker.storage.cli import (
|
|
||||||
_shm_period_and_invalid_count,
|
|
||||||
ldshm,
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def test_ldshm_detects_nulls_without_skewing_period() -> None:
|
|
||||||
'''
|
|
||||||
Null SHM slots must stop persistence without skewing cadence.
|
|
||||||
|
|
||||||
A live QQQ buffer contained an interior run of zero-time rows.
|
|
||||||
``ldshm`` included them in period inference and then sent them to
|
|
||||||
NativeDB, which correctly rejected the replacement. Build a
|
|
||||||
60-second series with the same interior hole and enough
|
|
||||||
bars around the gap. Prove inspection counts every null while
|
|
||||||
inferring the published series' 60-second period; the
|
|
||||||
command uses that count to skip both persistence and SHM reload.
|
|
||||||
|
|
||||||
'''
|
|
||||||
times = np.array([
|
|
||||||
60,
|
|
||||||
120,
|
|
||||||
0,
|
|
||||||
0,
|
|
||||||
300,
|
|
||||||
360,
|
|
||||||
420,
|
|
||||||
])
|
|
||||||
|
|
||||||
period_s, invalid_count = _shm_period_and_invalid_count(times)
|
|
||||||
|
|
||||||
assert period_s == 60
|
|
||||||
assert invalid_count == 2
|
|
||||||
|
|
||||||
|
|
||||||
def test_ldshm_period_is_an_observed_step() -> None:
|
|
||||||
'''
|
|
||||||
Sparse timestamps must not invent a storage timeframe.
|
|
||||||
|
|
||||||
The arithmetic median of two steps can produce a cadence that
|
|
||||||
occurs nowhere in the source series. Such a value can select the
|
|
||||||
wrong NativeDB timeframe and hide the gap. Arrange tied 60 and
|
|
||||||
120-second steps and prove the deterministic tie-break chooses
|
|
||||||
the smaller observed period.
|
|
||||||
|
|
||||||
'''
|
|
||||||
times = np.array([60, 120, 240])
|
|
||||||
|
|
||||||
period_s, invalid_count = _shm_period_and_invalid_count(times)
|
|
||||||
|
|
||||||
assert period_s == 60
|
|
||||||
assert invalid_count == 0
|
|
||||||
|
|
||||||
|
|
||||||
def test_ldshm_counts_every_invalid_timestamp() -> None:
|
|
||||||
'''
|
|
||||||
Corrupt timestamps stop persistence just like null SHM slots.
|
|
||||||
|
|
||||||
Exact zero denotes an unpublished slot, while negative, NaN, and
|
|
||||||
infinite values indicate corruption. NativeDB rejects them all,
|
|
||||||
but allowing a replacement attempt would still abort ``ldshm``.
|
|
||||||
Mix each invalid class into a regular series and prove the same
|
|
||||||
snapshot inspection counts all four while preserving cadence.
|
|
||||||
|
|
||||||
'''
|
|
||||||
times = np.array([
|
|
||||||
60,
|
|
||||||
120,
|
|
||||||
0,
|
|
||||||
-1,
|
|
||||||
np.nan,
|
|
||||||
np.inf,
|
|
||||||
180,
|
|
||||||
])
|
|
||||||
|
|
||||||
period_s, invalid_count = _shm_period_and_invalid_count(times)
|
|
||||||
|
|
||||||
assert period_s == 60
|
|
||||||
assert invalid_count == 4
|
|
||||||
|
|
||||||
|
|
||||||
def test_ldshm_invalid_snapshot_never_reaches_storage(
|
|
||||||
monkeypatch: pytest.MonkeyPatch,
|
|
||||||
) -> None:
|
|
||||||
'''
|
|
||||||
Invalid snapshots stop before analyzers and storage mutation.
|
|
||||||
|
|
||||||
The original QQQ failure passed zero-time rows through dedupe and
|
|
||||||
into NativeDB. An intermediate fix called null-segment analysis
|
|
||||||
before its guard, allowing corrupt indexes to abort the command.
|
|
||||||
Supply one immutable Polars snapshot containing that combination,
|
|
||||||
then replace dedupe and writes with failing spies. Completing the
|
|
||||||
command proves its guard runs first, preventing persistence or
|
|
||||||
SHM reload without relying on task timing.
|
|
||||||
|
|
||||||
'''
|
|
||||||
shm_df = pl.DataFrame({
|
|
||||||
'index': [100, 500, 104],
|
|
||||||
'time': [60, 0, 180],
|
|
||||||
})
|
|
||||||
|
|
||||||
@asynccontextmanager
|
|
||||||
async def open_runtime(
|
|
||||||
*args: Any,
|
|
||||||
**kwargs: Any,
|
|
||||||
) -> AsyncIterator[None]:
|
|
||||||
'''
|
|
||||||
Yield a runtime-free command context.
|
|
||||||
|
|
||||||
'''
|
|
||||||
yield
|
|
||||||
|
|
||||||
async def fail_write(
|
|
||||||
*args: Any,
|
|
||||||
**kwargs: Any,
|
|
||||||
) -> None:
|
|
||||||
'''
|
|
||||||
Fail if invalid SHM reaches NativeDB.
|
|
||||||
|
|
||||||
'''
|
|
||||||
raise AssertionError('invalid SHM reached NativeDB')
|
|
||||||
|
|
||||||
client = SimpleNamespace(write_ohlcv=fail_write)
|
|
||||||
|
|
||||||
@asynccontextmanager
|
|
||||||
async def open_storage(
|
|
||||||
*args: Any,
|
|
||||||
**kwargs: Any,
|
|
||||||
) -> AsyncIterator[tuple[SimpleNamespace, SimpleNamespace]]:
|
|
||||||
'''
|
|
||||||
Yield storage objects with a failing writer.
|
|
||||||
|
|
||||||
'''
|
|
||||||
yield SimpleNamespace(), client
|
|
||||||
|
|
||||||
@asynccontextmanager
|
|
||||||
async def open_annotations(
|
|
||||||
*args: Any,
|
|
||||||
**kwargs: Any,
|
|
||||||
) -> AsyncIterator[SimpleNamespace]:
|
|
||||||
'''
|
|
||||||
Yield an inert annotation controller.
|
|
||||||
|
|
||||||
'''
|
|
||||||
yield SimpleNamespace()
|
|
||||||
|
|
||||||
def iter_shms(
|
|
||||||
fqme: str,
|
|
||||||
) -> Iterator[tuple[Path, SimpleNamespace, pl.DataFrame]]:
|
|
||||||
'''
|
|
||||||
Yield the single captured invalid snapshot.
|
|
||||||
|
|
||||||
'''
|
|
||||||
yield Path('qqq.rt'), SimpleNamespace(), shm_df
|
|
||||||
|
|
||||||
def dedupe(*args: Any, **kwargs: Any) -> None:
|
|
||||||
'''
|
|
||||||
Fail if invalid SHM reaches deduplication.
|
|
||||||
|
|
||||||
'''
|
|
||||||
raise AssertionError('invalid SHM reached dedupe')
|
|
||||||
|
|
||||||
async def pause() -> None:
|
|
||||||
'''
|
|
||||||
Replace the command's final interactive pause.
|
|
||||||
|
|
||||||
'''
|
|
||||||
|
|
||||||
monkeypatch.setattr(
|
|
||||||
storage_cli,
|
|
||||||
'open_piker_runtime',
|
|
||||||
open_runtime,
|
|
||||||
)
|
|
||||||
monkeypatch.setattr(
|
|
||||||
storage_cli,
|
|
||||||
'open_storage_client',
|
|
||||||
open_storage,
|
|
||||||
)
|
|
||||||
monkeypatch.setattr(tsp, 'iter_dfs_from_shms', iter_shms)
|
|
||||||
monkeypatch.setattr(tsp, 'dedupe_ohlcv_smart', dedupe)
|
|
||||||
monkeypatch.setattr(storage_cli.tractor, 'pause', pause)
|
|
||||||
|
|
||||||
from piker.ui import _remote_ctl
|
|
||||||
monkeypatch.setattr(
|
|
||||||
_remote_ctl,
|
|
||||||
'open_annot_ctl',
|
|
||||||
open_annotations,
|
|
||||||
)
|
|
||||||
|
|
||||||
ldshm('qqq.nasdaq.ib')
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.parametrize(
|
|
||||||
'times',
|
|
||||||
[
|
|
||||||
np.array([0, 0]),
|
|
||||||
np.array([60]),
|
|
||||||
np.array([60, 60]),
|
|
||||||
],
|
|
||||||
)
|
|
||||||
def test_ldshm_skips_frame_without_positive_cadence(
|
|
||||||
times: np.ndarray,
|
|
||||||
) -> None:
|
|
||||||
'''
|
|
||||||
Degenerate buffers are skipped before indexing or persistence.
|
|
||||||
|
|
||||||
The old inference indexed final timestamps and left its median
|
|
||||||
undefined for short buffers. Fully unpublished, one-row, and
|
|
||||||
duplicate-only frames could therefore fail before the command's
|
|
||||||
guard. Exercise each shape and prove inspection
|
|
||||||
returns no cadence, which directs ``ldshm`` to skip the buffer.
|
|
||||||
|
|
||||||
'''
|
|
||||||
period_s, _ = _shm_period_and_invalid_count(times)
|
|
||||||
|
|
||||||
assert period_s is None
|
|
||||||
|
|
@ -129,69 +129,6 @@ def test_audit_preserves_raw_defect_evidence() -> None:
|
||||||
assert report['values']['infinity_by_column']['volume'] == 1
|
assert report['values']['infinity_by_column']['volume'] == 1
|
||||||
|
|
||||||
|
|
||||||
def test_audit_rejects_epoch_zero_price_sentinel() -> None:
|
|
||||||
'''
|
|
||||||
The legacy first-slot sentinel is structural corruption.
|
|
||||||
|
|
||||||
The repaired MNQ Parquet retained a row at epoch 60 whose OHLCV
|
|
||||||
payload was zero. Numeric and finite checks marked it valid even
|
|
||||||
though chart auto-ranging then included zero. Audit a canonical
|
|
||||||
frame with that sentinel and prove the report counts it and fails
|
|
||||||
structural qualification.
|
|
||||||
|
|
||||||
'''
|
|
||||||
frame = mk_frame((60, 120)).with_columns(
|
|
||||||
pl.Series('open', [0, 1], dtype=pl.Float64),
|
|
||||||
pl.Series('high', [0, 2], dtype=pl.Float64),
|
|
||||||
pl.Series('low', [0, 1], dtype=pl.Float64),
|
|
||||||
pl.Series('close', [0, 2], dtype=pl.Float64),
|
|
||||||
)
|
|
||||||
|
|
||||||
report: dict = audit_ohlcv_frame(
|
|
||||||
frame,
|
|
||||||
fqme='x.test',
|
|
||||||
period_s=60,
|
|
||||||
)
|
|
||||||
|
|
||||||
assert report['values']['all_zero_price_rows'] == 1
|
|
||||||
assert report['values']['epoch_sentinel_rows'] == 1
|
|
||||||
assert 'epoch_ohlc_sentinel' in report['result']['violations']
|
|
||||||
assert report['result']['structural_ok'] is False
|
|
||||||
|
|
||||||
|
|
||||||
def test_audit_leaves_modern_zero_prices_unclassified() -> None:
|
|
||||||
'''
|
|
||||||
Audit must not call every all-zero price bar corrupt.
|
|
||||||
|
|
||||||
A spread or synthetic instrument may legitimately trade at zero.
|
|
||||||
The known corruption is tied to the first slot. Place an all-zero
|
|
||||||
bar later and prove it remains visible as a warning without
|
|
||||||
failing structural qualification.
|
|
||||||
|
|
||||||
'''
|
|
||||||
frame = mk_frame((120,)).with_columns(
|
|
||||||
pl.Series('open', [0], dtype=pl.Float64),
|
|
||||||
pl.Series('high', [0], dtype=pl.Float64),
|
|
||||||
pl.Series('low', [0], dtype=pl.Float64),
|
|
||||||
pl.Series('close', [0], dtype=pl.Float64),
|
|
||||||
)
|
|
||||||
|
|
||||||
report: dict = audit_ohlcv_frame(
|
|
||||||
frame,
|
|
||||||
fqme='x.test',
|
|
||||||
period_s=60,
|
|
||||||
)
|
|
||||||
|
|
||||||
assert report['values']['all_zero_price_rows'] == 1
|
|
||||||
assert report['values']['epoch_sentinel_rows'] == 0
|
|
||||||
assert report['values']['unclassified_zero_price_rows'] == 1
|
|
||||||
assert report['result']['structural_ok'] is True
|
|
||||||
assert (
|
|
||||||
'all_zero_ohlc_prices_unclassified'
|
|
||||||
in report['result']['warnings']
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def test_gap_aggregates_are_not_truncated_with_details() -> None:
|
def test_gap_aggregates_are_not_truncated_with_details() -> None:
|
||||||
'''
|
'''
|
||||||
Limiting JSON detail must not undercount total missing coverage.
|
Limiting JSON detail must not undercount total missing coverage.
|
||||||
|
|
|
||||||
|
|
@ -131,162 +131,47 @@ def test_update_preserves_history_and_resolves_conflicts(
|
||||||
assert stored['index'].to_list() == [0, 1, 2, 3]
|
assert stored['index'].to_list() == [0, 1, 2, 3]
|
||||||
|
|
||||||
|
|
||||||
def test_update_repairs_invalid_persisted_sentinel_rows(
|
def test_update_repairs_nonpositive_persisted_timestamps(
|
||||||
tmp_path: Path,
|
tmp_path: Path,
|
||||||
) -> None:
|
) -> None:
|
||||||
'''
|
'''
|
||||||
Valid incoming history must repair legacy sentinel rows.
|
Valid incoming history must repair a legacy epoch-zero row.
|
||||||
|
|
||||||
The known MNQ baseline contained epoch-zero and all-zero rows,
|
The known MNQ baseline contains one zero timestamp plus extra derived
|
||||||
plus extra columns and absolute indexes. ``update_ohlcv()`` first
|
columns and noncanonical absolute indexes. ``update_ohlcv()`` used to
|
||||||
retained both, then removed only epoch zero and left chart range
|
canonicalize those bytes but retain the zero row, then reject the
|
||||||
anchored at zero. Write both sentinels, append a valid provider
|
whole merged frame before fresh IB bars could publish. Write that
|
||||||
frame, and prove only valid bars survive with canonical
|
legacy shape directly, append a valid provider frame, and prove only
|
||||||
indexes.
|
the invalid persisted row disappears while old and new valid bars
|
||||||
|
survive with canonical indexes.
|
||||||
|
|
||||||
'''
|
'''
|
||||||
client = NativeStorageClient(tmp_path)
|
client = NativeStorageClient(tmp_path)
|
||||||
legacy = (
|
legacy = (
|
||||||
tsp.np2pl(mk_ohlcv(
|
tsp.np2pl(mk_ohlcv(
|
||||||
(0, 60, 120, 180),
|
(0, 60, 120),
|
||||||
(0, 0, 2, 3),
|
(0, 1, 2),
|
||||||
))
|
))
|
||||||
.with_columns(
|
.with_columns(
|
||||||
pl.Series(
|
pl.Series('index', [3137800, 3137801, 3137802]),
|
||||||
'index',
|
pl.Series('time_prev', [None, 0, 60]),
|
||||||
[3137800, 3137801, 3137802, 3137803],
|
|
||||||
),
|
|
||||||
pl.Series('time_prev', [None, 0, 60, 120]),
|
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
path: Path = client.mk_path('x.test', 60)
|
path: Path = client.mk_path('x.test', 60)
|
||||||
legacy.write_parquet(path)
|
legacy.write_parquet(path)
|
||||||
|
|
||||||
loaded = trio.run(client.read_ohlcv, 'x.test', 60)
|
|
||||||
assert loaded['time'].tolist() == [120, 180]
|
|
||||||
assert loaded['index'].tolist() == [0, 1]
|
|
||||||
assert pl.read_parquet(path).height == 4
|
|
||||||
|
|
||||||
run(client.update_ohlcv(
|
run(client.update_ohlcv(
|
||||||
'x.test',
|
'x.test',
|
||||||
mk_ohlcv((240,), (4,)),
|
mk_ohlcv((180,), (3,)),
|
||||||
60,
|
60,
|
||||||
))
|
))
|
||||||
stored: pl.DataFrame = pl.read_parquet(path)
|
stored: pl.DataFrame = pl.read_parquet(path)
|
||||||
|
|
||||||
assert stored['time'].to_list() == [120, 180, 240]
|
assert stored['time'].to_list() == [60, 120, 180]
|
||||||
assert stored['close'].to_list() == [2, 3, 4]
|
assert stored['close'].to_list() == [1, 2, 3]
|
||||||
assert stored['index'].to_list() == [0, 1, 2]
|
assert stored['index'].to_list() == [0, 1, 2]
|
||||||
|
|
||||||
|
|
||||||
def test_write_preserves_modern_all_zero_prices(
|
|
||||||
tmp_path: Path,
|
|
||||||
) -> None:
|
|
||||||
'''
|
|
||||||
Zero prices outside the legacy sentinel range remain valid.
|
|
||||||
|
|
||||||
Some spreads or synthetic instruments may legitimately trade at
|
|
||||||
zero. Repair targets the observed first-slot sentinel, not every
|
|
||||||
all-zero OHLC bar. Write a modern all-zero row and
|
|
||||||
prove it survives normal replacement unchanged.
|
|
||||||
|
|
||||||
'''
|
|
||||||
client = NativeStorageClient(tmp_path)
|
|
||||||
zero_bar = mk_ohlcv((1_800_000_000,), (0,))
|
|
||||||
|
|
||||||
run(client.write_ohlcv('x.test', zero_bar, 60))
|
|
||||||
stored = pl.read_parquet(client.mk_path('x.test', 60))
|
|
||||||
|
|
||||||
assert stored['time'].to_list() == [1_800_000_000]
|
|
||||||
assert stored['close'].to_list() == [0]
|
|
||||||
|
|
||||||
|
|
||||||
def test_write_rejects_first_slot_zero_sentinel(
|
|
||||||
tmp_path: Path,
|
|
||||||
) -> None:
|
|
||||||
'''
|
|
||||||
New writes must not recreate the repaired legacy sentinel.
|
|
||||||
|
|
||||||
Read repair removes all-zero OHLC first-slot rows. Accepting one from
|
|
||||||
new input would recreate a file whose audit fails and whose read view
|
|
||||||
hides a row. Submit the exact 60-second
|
|
||||||
sentinel and prove no file publishes.
|
|
||||||
|
|
||||||
'''
|
|
||||||
client = NativeStorageClient(tmp_path)
|
|
||||||
sentinel = mk_ohlcv((60,), (0,))
|
|
||||||
|
|
||||||
with pytest.raises(ValueError, match='epoch sentinel'):
|
|
||||||
run(client.write_ohlcv('x.test', sentinel, 60))
|
|
||||||
|
|
||||||
assert not client.mk_path('x.test', 60).exists()
|
|
||||||
|
|
||||||
|
|
||||||
def test_sentinel_only_storage_loads_as_no_history(
|
|
||||||
tmp_path: Path,
|
|
||||||
) -> None:
|
|
||||||
'''
|
|
||||||
A fully repairable legacy file must permit fresh startup.
|
|
||||||
|
|
||||||
Filtering a file containing only sentinels yields an empty array.
|
|
||||||
``load()`` previously indexed its endpoint timestamps and crashed
|
|
||||||
before valid data could replace it. Persist only sentinel rows,
|
|
||||||
prove load reports no history, then append a valid bar and
|
|
||||||
verify durable storage becomes canonical.
|
|
||||||
|
|
||||||
'''
|
|
||||||
client = NativeStorageClient(tmp_path)
|
|
||||||
path: Path = client.mk_path('x.test', 60)
|
|
||||||
tsp.np2pl(mk_ohlcv(
|
|
||||||
(0, 60),
|
|
||||||
(0, 0),
|
|
||||||
)).write_parquet(path)
|
|
||||||
|
|
||||||
assert trio.run(client.load, 'x.test', 60) is None
|
|
||||||
run(client.update_ohlcv(
|
|
||||||
'x.test',
|
|
||||||
mk_ohlcv((120,), (1,)),
|
|
||||||
60,
|
|
||||||
))
|
|
||||||
stored: pl.DataFrame = pl.read_parquet(path)
|
|
||||||
|
|
||||||
assert stored['time'].to_list() == [120]
|
|
||||||
assert stored['index'].to_list() == [0]
|
|
||||||
|
|
||||||
|
|
||||||
def test_merge_does_not_hide_null_price_rows(
|
|
||||||
tmp_path: Path,
|
|
||||||
) -> None:
|
|
||||||
'''
|
|
||||||
Sentinel repair must not silently discard other malformed rows.
|
|
||||||
|
|
||||||
Polars boolean filters drop null predicates. Without an explicit
|
|
||||||
false fill, one null plus three zeros looked neither valid nor
|
|
||||||
all-zero but disappeared before validation. Persist that shape,
|
|
||||||
append valid data, and prove merge still rejects the null
|
|
||||||
while preserving the original evidence file.
|
|
||||||
|
|
||||||
'''
|
|
||||||
client = NativeStorageClient(tmp_path)
|
|
||||||
path: Path = client.mk_path('x.test', 60)
|
|
||||||
malformed = tsp.np2pl(
|
|
||||||
mk_ohlcv((1_800_000_000,), (0,))
|
|
||||||
).with_columns(
|
|
||||||
pl.Series('open', [None], dtype=pl.Float64)
|
|
||||||
)
|
|
||||||
malformed.write_parquet(path)
|
|
||||||
before: bytes = path.read_bytes()
|
|
||||||
|
|
||||||
with pytest.raises(ValueError, match='finite numeric'):
|
|
||||||
run(client.update_ohlcv(
|
|
||||||
'x.test',
|
|
||||||
mk_ohlcv((1_800_000_060,), (1,)),
|
|
||||||
60,
|
|
||||||
))
|
|
||||||
|
|
||||||
assert path.read_bytes() == before
|
|
||||||
|
|
||||||
|
|
||||||
def test_write_ohlcv_remains_an_explicit_replacement(
|
def test_write_ohlcv_remains_an_explicit_replacement(
|
||||||
tmp_path: Path,
|
tmp_path: Path,
|
||||||
) -> None:
|
) -> None:
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue