Files
fxhnt/tests/integration/test_orchestration_definitions.py
Jeroen Grusewski 04305a1c72 feat(orchestration): fold bybit precompute into a Dagster asset (K8sRunLauncher-ready)
Move _persist_bybit_book + _BYBIT_BOOK_SIDS to application/bybit_book_persist.py
(single source of truth; breaks the cli<->assets coupling). Add a new nightly
asset bybit_book_precompute (deps=[bybit_warehouse_refresh]) that wraps the
same helper backtest-refs/bybit-persist-sleeve-ret run — writing bybit_sleeve_ret
+ the report-kind gate refs in one compute-once pass. Tagged
dagster-k8s/config 6Gi req / 8Gi limit so K8sRunLauncher runs it in its own
right-sized Job (never the 4Gi daemon). Wired into combined_book_job + defs;
tests updated (13->14 assets + presence + wiring). CLI commands left intact
(deleted at the cron layer later). 60 tests pass.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-19 09:36:05 +00:00

141 lines
9.9 KiB
Python

"""The Definitions load and expose a daily schedule over the combined-book asset job + the 4 assets."""
from __future__ import annotations
def test_definitions_load_with_assets_and_schedule() -> None:
from fxhnt.adapters.orchestration.definitions import defs
# --- assets present ---
# defs.get_repository_def() is available in Dagster 1.12; its .assets_defs_by_key
# maps AssetKey → AssetsDefinition for every registered asset.
repo = defs.get_repository_def()
asset_keys = {k.to_user_string() for k in repo.assets_defs_by_key}
assert {"futures_bars", "cockpit_forward"} <= asset_keys
# Phase 0b venue consolidation (Task 7b): the Binance crypto_bars/crypto_funding/crypto_spot ingest
# assets are retired — crypto is 100% Bybit.
for retired_binance_ingest in ("crypto_bars", "crypto_funding", "crypto_spot"):
assert retired_binance_ingest not in asset_keys, f"{retired_binance_ingest} must be retired"
# B1: the surviving paper-track assets (sixtyforty kept as the 60/40 benchmark).
# Retired: crossvenue (falsified 3x), funding_nav (superseded by xsfunding), gd_nav (futures premia —
# falsified, scale-gated) + poc_nav (old crypto PoC, superseded). NOTE: multistrat_nav was RE-ADDED
# 2026-07-08 (TradFi diversifier on Alpaca, uncorrelated to crypto) — no longer retired.
assert "sixtyforty_nav" in asset_keys, f"missing sixtyforty benchmark: {asset_keys}"
for retired in ("crossvenue_nav", "funding_nav", "gd_nav", "poc_nav"):
assert retired not in asset_keys, f"{retired} must be retired"
# equity-factor sleeves RETIRED 2026-06-20 (weakly-held, edge never established)
for retired_eq in ("eqfactor_long_nav", "eqfactor_tilt_nav", "eqfactor_scores", "eqfactor_ls_nav"):
assert retired_eq not in asset_keys, f"{retired_eq} must be retired"
# Phase 0b Task 4 (2026-07-13): standalone Binance crypto_tstrend_nav + stablecoin_rotation_nav
# RETIRED (neither is tradeable on Bybit — the crypto_tstrend SLEEVE inside the bybit_4edge deploy
# book is a separate, unchanged thing).
for retired_binance in ("crypto_tstrend_nav", "stablecoin_rotation_nav"):
assert retired_binance not in asset_keys, f"{retired_binance} must be retired"
# the 2 validated crypto edges are wired in
assert "xsfunding_nav" in asset_keys, f"missing xsfunding_nav asset: {asset_keys}"
assert "unlock_nav" in asset_keys, f"missing unlock_nav asset: {asset_keys}"
# multistrat_nav (2026-07-08): TradFi multi-strat book, uncorrelated diversifier, on the reconciliation gate
assert "multistrat_nav" in asset_keys, f"missing multistrat_nav: {asset_keys}"
# vrp_nav ARCHIVED/FALSIFIED 2026-07-15 — the VRP sleeve is a shelved dead-end; its asset FUNCTION is
# KEPT in assets.py (so test_vrp_nav_asset still imports it), but it is DE-WIRED from the nightly, so it
# no longer registers as a Dagster asset. See project_fxhnt_diversifier_hunt_2026_07_15.
assert "vrp_nav" not in asset_keys, f"vrp_nav must be de-wired (archived): {asset_keys}"
# paper_book_snapshot (T7) RETIRED (Phase 0b Task 7d): the dead Binance-era combined-crypto `/paper`
# book — redundant with the `bybit_4edge` deploy book below.
assert "paper_book_snapshot" not in asset_keys, "paper_book_snapshot must be retired"
# Bybit forward paper track (out-of-sample record): nightly warehouse refresh + naive eq-wt 4-edge NAV.
assert "bybit_warehouse_refresh" in asset_keys, f"missing bybit_warehouse_refresh: {asset_keys}"
# bybit_book_precompute (2026-07-19): the heavy 4-sleeve precompute folded in from the retired
# compare-measured-precompute + backtest-refs CronJobs — writes bybit_sleeve_ret (cockpit-sim precompute)
# + the report-kind reconciliation-gate refs in one compute-once pass. Runs in an own-memory per-run K8s
# Job (tagged 6Gi/8Gi), depends on bybit_warehouse_refresh.
assert "bybit_book_precompute" in asset_keys, f"missing bybit_book_precompute: {asset_keys}"
assert "bybit_4edge_nav" in asset_keys, f"missing bybit_4edge_nav: {asset_keys}"
# observe-only levered shadow of the 4-edge book — independently recomputes from the warehouse, own state file
assert "bybit_4edge_levered_nav" in asset_keys, f"missing bybit_4edge_levered_nav: {asset_keys}"
# deribit_funding_bars: nightly Deribit perp funding ingest (the Deribit leg of the cross-venue carry)
assert "deribit_funding_bars" in asset_keys, f"missing deribit_funding_bars: {asset_keys}"
# positioning sleeve surfaced as its own forward track
assert "positioning_nav" in asset_keys, f"missing positioning_nav: {asset_keys}"
# xvenue_carry_nav RETIRED 2026-07-07 (validated-but-not-live, ~0 contribution); deribit ingest KEPT
assert "xvenue_carry_nav" not in asset_keys, "xvenue_carry_nav must be retired"
# Bybit LIVE paper book (per-symbol positions/trades/MTM, the deploy venue's live record) — depends on
# bybit_warehouse_refresh (READ-ONLY against it).
assert "bybit_paper_book" in asset_keys, f"missing bybit_paper_book: {asset_keys}"
# the Bybit testnet real-paper execution leg (formerly a separate job, Task 8) was REMOVED (Task 7):
# the leg is now CLI-driven via `run_bybit_testnet` / `execute-bybit`, not a Dagster asset.
assert "bybit_testnet_reconcile" not in asset_keys, "bybit_testnet_reconcile must be removed"
assert "bybit_testnet_record" not in asset_keys, "bybit_testnet_record must be removed"
# 14 assets: futures_bars, sixtyforty_nav,
# xsfunding_nav, unlock_nav, multistrat_nav, multistrat_levered_nav,
# deribit_funding_bars, bybit_warehouse_refresh, bybit_book_precompute, bybit_4edge_nav,
# bybit_4edge_levered_nav, positioning_nav, bybit_paper_book, cockpit_forward
# (xvenue_carry_nav retired 2026-07-07; multistrat_nav re-added 2026-07-08;
# bybit_testnet_reconcile/bybit_testnet_record removed 2026-07-11, Task 7;
# combined_forward_nav retired 2026-07-11 — vestigial crypto-momentum book;
# vrp_nav added 2026-07-12, Task 6;
# crypto_tstrend_nav/stablecoin_rotation_nav retired 2026-07-13, Phase 0b Task 4;
# crypto_bars/crypto_funding/crypto_spot Binance ingest retired 2026-07-13, Phase 0b Task 7b;
# paper_book_snapshot retired 2026-07-14, Phase 0b Task 7d — the dead combined-crypto book;
# vrp_nav ARCHIVED/de-wired 2026-07-15 — VRP falsified/shelved, asset fn kept but unregistered;
# multistrat_levered_nav added 2026-07-18 — observe-only levered shadow of the ETF book;
# bybit_book_precompute added 2026-07-19 — heavy 4-sleeve precompute folded in from the retired
# compare-measured-precompute + backtest-refs CronJobs)
assert len(asset_keys) == 14, f"expected 14 assets, got {len(asset_keys)}: {asset_keys}"
# bybit_paper_book is READ-ONLY against the warehouse: its ONLY upstream is bybit_warehouse_refresh.
from dagster import AssetKey as _AK
bpb = repo.assets_defs_by_key[_AK("bybit_paper_book")]
bpb_up = {k.to_user_string() for k in bpb.asset_deps[_AK("bybit_paper_book")]}
assert bpb_up == {"bybit_warehouse_refresh"}, f"bybit_paper_book upstream={bpb_up}"
# --- schedule present with the right cron ---
# defs.schedules is a list[ScheduleDefinition] (or None when empty)
scheds = list(defs.schedules or [])
assert any(getattr(s, "cron_schedule", "") == "30 23 * * *" for s in scheds), (
f"expected a schedule with cron '30 23 * * *'; got {[getattr(s, 'cron_schedule', None) for s in scheds]}"
)
# auto-arm on fresh deploy — else a recreated instance silently never runs the nightly track
from dagster import DefaultScheduleStatus
daily = next(s for s in scheds if getattr(s, "cron_schedule", "") == "30 23 * * *")
assert daily.default_status == DefaultScheduleStatus.RUNNING
# --- cockpit_forward depends on every paper-track asset (so it ingests their state files) ---
from dagster import AssetKey
cockpit_def = repo.assets_defs_by_key[AssetKey("cockpit_forward")]
deps = cockpit_def.asset_deps[AssetKey("cockpit_forward")]
upstream = {k.to_user_string() for k in deps}
expected_upstream = {
"xsfunding_nav", "unlock_nav", "sixtyforty_nav",
"multistrat_nav", "bybit_4edge_nav", "positioning_nav",
}
assert expected_upstream <= upstream, f"cockpit_forward missing upstream: {expected_upstream - upstream}"
# vrp_nav ARCHIVED/de-wired 2026-07-15 — no longer a cockpit_forward dep (nightly stops materializing it).
assert "vrp_nav" not in upstream, f"vrp_nav must be de-wired from cockpit_forward: {upstream}"
def test_bybit_book_precompute_wired_into_job_and_defs():
"""The heavy 4-sleeve precompute (folded in from the retired compare-measured-precompute + backtest-refs
CronJobs) must be BOTH a registered asset AND in the nightly `combined_book_forward_job` selection — else
it would exist but never materialize on the schedule."""
from fxhnt.adapters.orchestration import definitions as d
# in the Definitions asset set
asset_keys = {k.to_user_string() for k in d.defs.resolve_all_asset_keys()}
assert "bybit_book_precompute" in asset_keys, f"missing from defs.assets: {asset_keys}"
# in the combined_book_forward_job selection (so the daily schedule materializes it)
job_keys = {k.to_user_string() for k in d.combined_book_job.selection.resolve(list(d.defs.assets))}
assert "bybit_book_precompute" in job_keys, (
f"bybit_book_precompute must be in combined_book_forward_job selection: {job_keys}")
def test_bybit_testnet_exec_removed_from_definitions():
from fxhnt.adapters.orchestration.definitions import defs
job_names = {j.name for j in defs.jobs}
assert "bybit_testnet_execution_job" not in job_names
sched_names = {s.name for s in defs.schedules}
assert "bybit_testnet_daily" not in sched_names
asset_keys = {k.to_user_string() for k in defs.resolve_all_asset_keys()}
assert "bybit_testnet_reconcile" not in asset_keys
assert "bybit_testnet_record" not in asset_keys