run_bb_squeeze_t_gated_observer.py 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282
  1. from __future__ import annotations
  2. import argparse
  3. import json
  4. import sys
  5. import time
  6. from dataclasses import asdict, dataclass
  7. from datetime import UTC, datetime
  8. from pathlib import Path
  9. import pandas as pd
  10. sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
  11. from okx_codex_trader.models import Candle
  12. from okx_codex_trader.okx_client import OkxClient
  13. ROOT = Path(__file__).resolve().parents[1]
  14. STATE_DIR = ROOT / "var" / "bb-squeeze-t-gated-observer"
  15. STATE_FILE = "runtime-state.json"
  16. EVENTS_FILE = "observer-events.jsonl"
  17. ETH_SYMBOL = "ETH-USDT-SWAP"
  18. BTC_SYMBOL = "BTC-USDT-SWAP"
  19. BAR = "15m"
  20. LIVE_CANDLE_LIMIT = 1_200
  21. BAND_LENGTH = 96
  22. BANDWIDTH_LOOKBACK = 960
  23. BANDWIDTH_QUANTILE = 0.25
  24. STOP_LOSS_PCT = 0.01
  25. EXTREME_TAKE_PROFIT_PCT = 0.035
  26. ETH_VOL_CAP = 0.006
  27. COOLDOWN_BARS = 24
  28. REENTRY_BARS = 48
  29. BTC_TREND = 480
  30. BTC_MOMENTUM = 96
  31. MIN_FRAME_ROWS = BAND_LENGTH + BANDWIDTH_LOOKBACK
  32. @dataclass(frozen=True)
  33. class ObserverState:
  34. last_candle_ts: int | None
  35. active_side: str | None
  36. entry_price: float | None
  37. entry_candle_ts: int | None
  38. cooldown_until_ts: int | None
  39. reentry_side: str | None
  40. reentry_anchor_price: float | None
  41. reentry_until_ts: int | None
  42. EMPTY_STATE = ObserverState(None, None, None, None, None, None, None, None)
  43. def now_iso() -> str:
  44. return datetime.now(UTC).isoformat(timespec="seconds").replace("+00:00", "Z")
  45. def load_state(path: Path) -> ObserverState:
  46. if not path.exists():
  47. return EMPTY_STATE
  48. payload = json.loads(path.read_text(encoding="utf-8"))
  49. return ObserverState(
  50. last_candle_ts=payload["last_candle_ts"],
  51. active_side=payload["active_side"],
  52. entry_price=payload["entry_price"],
  53. entry_candle_ts=payload["entry_candle_ts"],
  54. cooldown_until_ts=payload["cooldown_until_ts"],
  55. reentry_side=payload["reentry_side"],
  56. reentry_anchor_price=payload["reentry_anchor_price"],
  57. reentry_until_ts=payload["reentry_until_ts"],
  58. )
  59. def save_state(path: Path, state: ObserverState) -> None:
  60. path.parent.mkdir(parents=True, exist_ok=True)
  61. path.write_text(json.dumps(asdict(state), indent=2, sort_keys=True) + "\n", encoding="utf-8")
  62. def append_jsonl(path: Path, payload: dict[str, object]) -> None:
  63. path.parent.mkdir(parents=True, exist_ok=True)
  64. with path.open("a", encoding="utf-8") as handle:
  65. handle.write(json.dumps(payload, sort_keys=True, separators=(",", ":")) + "\n")
  66. def frame_from_candles(candles: list[Candle]) -> pd.DataFrame:
  67. frame = pd.DataFrame([asdict(candle) for candle in candles])
  68. frame["time"] = pd.to_datetime(frame["ts"], unit="ms", utc=True)
  69. return frame.sort_values("ts").drop_duplicates("ts", keep="last").reset_index(drop=True)
  70. def load_pair_frame(client: OkxClient) -> pd.DataFrame:
  71. eth = frame_from_candles(client.get_candles(ETH_SYMBOL, BAR, LIVE_CANDLE_LIMIT))
  72. btc = frame_from_candles(client.get_candles(BTC_SYMBOL, BAR, LIVE_CANDLE_LIMIT))
  73. btc = btc[["ts", "close"]].rename(columns={"close": "btc_close"})
  74. return eth.merge(btc, on="ts", how="inner").sort_values("ts").reset_index(drop=True)
  75. def signal_from_frame(frame: pd.DataFrame, state: ObserverState) -> tuple[ObserverState, dict[str, object]]:
  76. if len(frame) < MIN_FRAME_ROWS:
  77. raise ValueError("not enough candles")
  78. eth_close = frame["close"].astype(float)
  79. btc_close = frame["btc_close"].astype(float)
  80. middle = eth_close.rolling(BAND_LENGTH).mean()
  81. stdev = eth_close.rolling(BAND_LENGTH).std(ddof=0)
  82. upper = middle + 2.0 * stdev
  83. lower = middle - 2.0 * stdev
  84. bandwidth = (upper - lower) / middle
  85. threshold = bandwidth.rolling(BANDWIDTH_LOOKBACK).quantile(BANDWIDTH_QUANTILE)
  86. eth_vol = eth_close.pct_change().rolling(96).std(ddof=0)
  87. btc_sma = btc_close.rolling(BTC_TREND).mean()
  88. btc_momentum = btc_close / btc_close.shift(BTC_MOMENTUM) - 1.0
  89. index = len(frame) - 1
  90. row = frame.iloc[index]
  91. candle_ts = int(row["ts"])
  92. candle_time = pd.Timestamp(row["time"]).isoformat().replace("+00:00", "Z")
  93. indicators = {
  94. "eth_close": float(row["close"]),
  95. "btc_close": float(row["btc_close"]),
  96. "middle": float(middle.iloc[index]),
  97. "upper": float(upper.iloc[index]),
  98. "lower": float(lower.iloc[index]),
  99. "bandwidth": float(bandwidth.iloc[index]),
  100. "bandwidth_threshold": float(threshold.iloc[index]),
  101. "eth_vol_96": float(eth_vol.iloc[index]),
  102. "btc_sma_480": float(btc_sma.iloc[index]),
  103. "btc_momentum_96": float(btc_momentum.iloc[index]),
  104. }
  105. if state.last_candle_ts is not None and candle_ts <= state.last_candle_ts:
  106. return state, {
  107. "decision_candle_ts": candle_ts,
  108. "decision_candle_time": candle_time,
  109. "signal": "state_replay",
  110. "target_side": state.active_side or "flat",
  111. "indicators": indicators,
  112. }
  113. next_state = ObserverState(
  114. candle_ts,
  115. state.active_side,
  116. state.entry_price,
  117. state.entry_candle_ts,
  118. state.cooldown_until_ts,
  119. state.reentry_side,
  120. state.reentry_anchor_price,
  121. state.reentry_until_ts,
  122. )
  123. signal = "hold"
  124. target_side = state.active_side or "flat"
  125. reentry_gate = False
  126. if state.active_side is not None:
  127. entry_price = float(state.entry_price)
  128. stop = entry_price * (1.0 - STOP_LOSS_PCT if state.active_side == "long" else 1.0 + STOP_LOSS_PCT)
  129. extreme_take = entry_price * (1.0 + EXTREME_TAKE_PROFIT_PCT if state.active_side == "long" else 1.0 - EXTREME_TAKE_PROFIT_PCT)
  130. stop_hit = (state.active_side == "long" and float(row["low"]) <= stop) or (state.active_side == "short" and float(row["high"]) >= stop)
  131. extreme_take_hit = (state.active_side == "long" and float(row["high"]) >= extreme_take) or (
  132. state.active_side == "short" and float(row["low"]) <= extreme_take
  133. )
  134. middle_exit = (state.active_side == "long" and float(row["close"]) < indicators["middle"]) or (
  135. state.active_side == "short" and float(row["close"]) > indicators["middle"]
  136. )
  137. if stop_hit or extreme_take_hit or middle_exit:
  138. signal = "exit_stop" if stop_hit else "exit_extreme_take" if extreme_take_hit else "exit_middle"
  139. target_side = "flat"
  140. if extreme_take_hit and not stop_hit:
  141. next_state = ObserverState(
  142. candle_ts,
  143. None,
  144. None,
  145. None,
  146. state.cooldown_until_ts,
  147. state.active_side,
  148. extreme_take,
  149. candle_ts + REENTRY_BARS * 900_000,
  150. )
  151. else:
  152. next_state = ObserverState(candle_ts, None, None, None, candle_ts + COOLDOWN_BARS * 900_000, None, None, None)
  153. else:
  154. if state.reentry_side is not None:
  155. if state.reentry_until_ts is None or candle_ts > state.reentry_until_ts:
  156. next_state = ObserverState(candle_ts, None, None, None, state.cooldown_until_ts, None, None, None)
  157. else:
  158. reentry_gate = (state.reentry_side == "long" and indicators["btc_momentum_96"] < 0.0) or (
  159. state.reentry_side == "short" and indicators["btc_momentum_96"] > 0.0
  160. )
  161. if reentry_gate:
  162. signal = "reentry_" + state.reentry_side
  163. target_side = state.reentry_side
  164. next_state = ObserverState(candle_ts, state.reentry_side, float(row["close"]), candle_ts, state.cooldown_until_ts, None, None, None)
  165. else:
  166. cooldown_ok = state.cooldown_until_ts is None or candle_ts >= state.cooldown_until_ts
  167. compressed = indicators["bandwidth"] <= indicators["bandwidth_threshold"]
  168. vol_ok = indicators["eth_vol_96"] <= ETH_VOL_CAP
  169. btc_up = indicators["btc_close"] > indicators["btc_sma_480"]
  170. if cooldown_ok and compressed and vol_ok and btc_up and float(row["close"]) > indicators["upper"]:
  171. signal = "entry_long"
  172. target_side = "long"
  173. next_state = ObserverState(candle_ts, "long", float(row["close"]), candle_ts, state.cooldown_until_ts, None, None, None)
  174. elif cooldown_ok and compressed and vol_ok and btc_up and float(row["close"]) < indicators["lower"]:
  175. signal = "entry_short"
  176. target_side = "short"
  177. next_state = ObserverState(candle_ts, "short", float(row["close"]), candle_ts, state.cooldown_until_ts, None, None, None)
  178. return next_state, {
  179. "decision_candle_ts": candle_ts,
  180. "decision_candle_time": candle_time,
  181. "signal": signal,
  182. "target_side": target_side,
  183. "reentry_gate": reentry_gate,
  184. "indicators": indicators,
  185. "params": {
  186. "band_length": BAND_LENGTH,
  187. "bandwidth_lookback": BANDWIDTH_LOOKBACK,
  188. "bandwidth_quantile": BANDWIDTH_QUANTILE,
  189. "stop_loss_pct": STOP_LOSS_PCT,
  190. "extreme_take_profit_pct": EXTREME_TAKE_PROFIT_PCT,
  191. "eth_vol_cap": ETH_VOL_CAP,
  192. "cooldown_bars": COOLDOWN_BARS,
  193. "reentry_bars": REENTRY_BARS,
  194. "entry_btc_filter": "btc-up",
  195. "reentry_gate_mode": "btc_against",
  196. },
  197. }
  198. def run_once(state_dir: Path) -> dict[str, object]:
  199. state_dir.mkdir(parents=True, exist_ok=True)
  200. state_path = state_dir / STATE_FILE
  201. previous_state = load_state(state_path)
  202. frame = load_pair_frame(OkxClient())
  203. next_state, signal = signal_from_frame(frame, previous_state)
  204. save_state(state_path, next_state)
  205. payload = {
  206. "created_at": now_iso(),
  207. "mode": "bb_squeeze_t_gated_readonly_observer",
  208. "orders_submitted": 0,
  209. "strategy": "bb-squeeze-t-l96-bw960-q0.25-sl0.01-xtp0.035-both-btc-up-vc0.006-dd0.25-cd24-tre48-btc_against",
  210. "previous_state": asdict(previous_state),
  211. "next_state": asdict(next_state),
  212. "signal": signal,
  213. "candles": {
  214. "rows": len(frame),
  215. "first_time": pd.Timestamp(frame.iloc[0]["time"]).isoformat().replace("+00:00", "Z"),
  216. "last_time": pd.Timestamp(frame.iloc[-1]["time"]).isoformat().replace("+00:00", "Z"),
  217. },
  218. "risk_limits": {
  219. "no_order_submission": True,
  220. "no_cancel_submission": True,
  221. "execution": "read_only_signal_stream",
  222. },
  223. }
  224. (state_dir / "heartbeat.json").write_text(json.dumps(payload, indent=2, sort_keys=True) + "\n", encoding="utf-8")
  225. append_jsonl(state_dir / EVENTS_FILE, payload)
  226. return payload
  227. def main() -> int:
  228. parser = argparse.ArgumentParser(description="Run BB squeeze T-gated read-only observer.")
  229. parser.add_argument("--state-dir", type=Path, default=STATE_DIR)
  230. parser.add_argument("--interval-seconds", type=int, default=300)
  231. parser.add_argument("--once", action="store_true")
  232. args = parser.parse_args()
  233. while True:
  234. try:
  235. payload = run_once(args.state_dir)
  236. print(json.dumps(payload, indent=2, sort_keys=True))
  237. except Exception as exc:
  238. error = {"created_at": now_iso(), "mode": "bb_squeeze_t_gated_readonly_observer", "orders_submitted": 0, "error": str(exc)}
  239. append_jsonl(args.state_dir / EVENTS_FILE, error)
  240. print(json.dumps(error, indent=2, sort_keys=True), file=sys.stderr)
  241. if args.once:
  242. return 1
  243. if args.once:
  244. return 0
  245. time.sleep(args.interval_seconds)
  246. if __name__ == "__main__":
  247. raise SystemExit(main())