| Overall Statistics |
|
Total Orders 3978 Average Win 0.94% Average Loss -0.57% Compounding Annual Return 20.513% Drawdown 30.800% Expectancy 0.110 Start Equity 100000 End Equity 306652.3 Net Profit 206.652% Sharpe Ratio 0.708 Sortino Ratio 0.961 Probabilistic Sharpe Ratio 13.842% Loss Rate 58% Win Rate 42% Profit-Loss Ratio 1.65 Alpha 0 Beta 0 Annual Standard Deviation 0.174 Annual Variance 0.03 Information Ratio 0.889 Tracking Error 0.174 Treynor Ratio 0 Total Fees $8552.70 Estimated Strategy Capacity $170000000.00 Lowest Capacity Asset NQ YYFADOG4CO3L Portfolio Turnover 308.10% Drawdown Recovery 238 |
# region imports
from AlgorithmImports import *
# endregion
# indicators/atr.py
from collections import deque
class ATR14:
"""
14-period Average True Range using a simple rolling mean.
Matches the original add_consolidation_indicators() exactly:
TR = max(H-L, |H-prev_C|, |L-prev_C|)
ATR = rolling mean of TR over 14 bars
"""
def __init__(self):
self._prev_close = None
self._tr_window = deque(maxlen=14)
self.value = None # None until 14 bars have accumulated
def update(self, bar) -> None:
high = float(bar.High)
low = float(bar.Low)
close = float(bar.Close)
if self._prev_close is None:
tr = high - low
else:
hl = high - low
hc = abs(high - self._prev_close)
lc = abs(low - self._prev_close)
tr = max(hl, hc, lc)
self._prev_close = close
self._tr_window.append(tr)
if len(self._tr_window) == 14:
self.value = sum(self._tr_window) / 14.0
def reset(self) -> None:
self._prev_close = None
self._tr_window.clear()
self.value = None
def is_ready(self) -> bool:
return self.value is not Nonefrom AlgorithmImports import *
def handle_15m_bar(algo, bar):
# Keep the mapped contract explicitly subscribed across rolls.
algo._execution.ensure_mapped_subscription()
# -- 1. Skip missing Observix bars -----------------------------------------
bar_key = bar.end_time.strftime("%Y-%m-%d %H:%M:%S")
if bar_key in algo._missing_bars:
return
# -- 1b. Inline daily reset (v2.5) -----------------------------------------
# Must fire BEFORE entry logic. The scheduled _daily_reset at 9:30 ET fires
# AFTER the bar handler for the same timestamp, causing first-of-day trades
# to inherit prior-day's seq counter (bug: first trade of 2026-04-23 tagged 003).
bar_date = bar.end_time.date()
if algo._session_date != bar_date:
algo._daily_trade_seq = 0
algo._session_pnl = 0.0
algo._session_trades = 0
algo._session_date = bar_date
try:
algo._day_start_equity = algo.portfolio.total_portfolio_value
except Exception:
pass
# -- 2. Warm-up gate -------------------------------------------------------
pre_trade = bar.end_time.replace(tzinfo=None) < algo._trade_start
# -- 3. Update all indicators (always, including warm-up) --------------------
algo.atr14.update(bar)
algo.chop_filter.update(algo.atr14.value)
algo._last_close = float(bar.Close)
algo.signal_engine.update_indicators(bar)
# -- 3b. Startup reconciliation (live-only, fires once on first real bar) --
# Detect desync between broker holdings and tracked _position_side at launch.
# If broker has a position we didn't track, synthesize a trade so exits work.
# If broker is flat but state claims position, reset to None.
if not algo._startup_reconciled and algo.live_mode:
algo._startup_reconciled = True
try:
holdings = algo.portfolio[algo.nq.mapped]
broker_qty = holdings.quantity if holdings else 0
except Exception as e:
algo.log(f"STARTUP RECONCILE: broker lookup error: {e}")
broker_qty = 0
if broker_qty > 0 and algo._position_side is None:
algo._position_side = "LONG"
algo.trade_tracker.open_trade("LONG", algo._last_close, None, None, algo.time)
algo.log(f"STARTUP RECONCILE: broker has {broker_qty}, set to LONG, "
f"synthetic trade at {algo._last_close}")
# algo.telegram.send_telegram_message(
# f"\u26a0\ufe0f Startup: detected untracked LONG ({broker_qty} contracts)"
# )
elif broker_qty < 0 and algo._position_side is None:
algo._position_side = "SHORT"
algo.trade_tracker.open_trade("SHORT", algo._last_close, None, None, algo.time)
algo.log(f"STARTUP RECONCILE: broker has {broker_qty}, set to SHORT, "
f"synthetic trade at {algo._last_close}")
# algo.telegram.send_telegram_message(
# f"\u26a0\ufe0f Startup: detected untracked SHORT ({broker_qty} contracts)"
# )
elif broker_qty == 0 and algo._position_side is not None:
algo.log(f"STARTUP RECONCILE: broker flat but state={algo._position_side}, resetting")
algo._position_side = None
algo._stop_loss = None
algo._take_profit = None
if algo.trade_tracker.is_open:
algo.trade_tracker.reset()
else:
algo.log(f"STARTUP RECONCILE: OK (broker={broker_qty}, state={algo._position_side})")
# -- 3c. Backtest position self-healing (contract expiry guard) ----------
# In backtest mode, on_order_event() returns immediately (live_mode=False),
# so QC's automatic liquidation at contract expiry never resets _position_side.
# Result: _position_side stays "LONG"/"SHORT" while portfolio is flat,
# blocking ALL new entries for the remainder of the backtest.
# Fix: detect the mismatch every bar and self-heal.
# This is safe in live mode too -- the startup reconcile + on_order_event
# already handle live drift, but an extra flat-check here costs nothing.
if algo._position_side is not None and not algo.portfolio.invested:
algo.log(
f"STATE SYNC: portfolio flat but _position_side={algo._position_side} "
f"at {bar.end_time} -- likely contract expiry auto-liquidation. Resetting."
)
if algo.trade_tracker.is_open:
# Close the tracker so records are saved; use last close as proxy price.
record = algo.trade_tracker.close_trade(
algo._last_close, bar.end_time, "Contract expiry / auto-liquidation"
)
algo._trade_log.append(record)
algo._position_side = None
algo._execution.clear_position()
algo._stop_loss = None
algo._take_profit = None
algo._active_contract_size = algo.config["trade"]["contract_size"] # v6: reset
# -- 4. Snapshot position state at bar open ------------------------------
was_invested = algo._position_side is not None
already_long = algo._position_side == "LONG"
already_short = algo._position_side == "SHORT"
exit_side = algo._position_side if was_invested else "LONG"
sl_level = algo._stop_loss
tp_level = algo._take_profit
# -- 5. Exit signal (always update, routed by SignalEngine) -----------------
exit_fired = algo.signal_engine.check_exit(bar, exit_side, bar.end_time)
# -- 6. Exit checks (only if invested at bar open, past trade start) ---
sl_hit = False
tp_hit = False
exit_price = None
exit_reason = ""
if was_invested and not pre_trade:
sl_hit, tp_hit, exit_price, exit_reason = \
algo.trade_tracker.update_bar(bar)
if exit_fired and not sl_hit and not tp_hit:
exit_reason = (f"{'Long' if algo._position_side == 'LONG' else 'Short'} "
f"exit signal.")
algo._close_position(float(bar.Close), bar.end_time, exit_reason)
elif sl_hit or tp_hit:
algo._close_position(exit_price, bar.end_time, exit_reason)
# -- 6b. Accumulate exit indicator row (only past trade start) ---------
ex = algo.signal_engine.last_exit
if not pre_trade and ex.get("ma7_c") is not None and ex.get("ma7_p") is not None:
algo._exit_rows.append({
"time": algo.time,
"close": algo._last_close,
"high": float(bar.High),
"low": float(bar.Low),
"otf_up": ex["otf_up"],
"otf_down": ex["otf_down"],
"ma7_c": ex["ma7_c"],
"ma7_p": ex["ma7_p"],
"flat": ex["flat"],
"direction": ex["direction"],
"exit": ex["exit"],
"exit_signal": ex.get("exit_signal"),
"sl_level": sl_level,
"tp_level": tp_level,
"sl_hit": sl_hit,
"tp_hit": tp_hit,
"was_invested": was_invested,
"exit_reason": exit_reason,
})
# -- 7. Entry signal (always update, routed by SignalEngine) ----------------
cfg_entry = algo.config["entry"]
(entry_long, entry_short,
sl_long, tp_long,
sl_short, tp_short) = algo.signal_engine.update_entry(
bar = bar,
atr14_value = algo.atr14.value,
chop_ok = algo.chop_filter.value,
timestamp = bar.end_time,
cfg_entry = cfg_entry,
)
# -- 8. Trade-state gating (only past trade start, only if flat) -------
rth_ok = algo._is_rth() if algo._rth_only else True
regime_ok = algo.hmm_regime.is_good_regime(bar.end_time) if algo.hmm_regime is not None else True
broker_flat = True
if algo.live_mode:
try:
h = algo.portfolio[algo.nq.mapped]
broker_flat = (h is None or h.quantity == 0)
except Exception:
broker_flat = True
if not pre_trade and algo._position_side is None and broker_flat and rth_ok and regime_ok:
_sl = sl_long if algo._sl_tp_enabled else None
_tp = tp_long if algo._sl_tp_enabled else None
_sl_s = sl_short if algo._sl_tp_enabled else None
_tp_s = tp_short if algo._sl_tp_enabled else None
if entry_long:
algo._open_position("LONG", float(bar.Close),
_sl, _tp, bar.end_time, bar)
elif entry_short and not algo._long_only:
algo._open_position("SHORT", float(bar.Close),
_sl_s, _tp_s, bar.end_time, bar)
# -- 9. Accumulate entry indicator row (only past trade start) ---------
en = algo.signal_engine.last_entry
if not pre_trade and en.get("ma7_c") is not None and en.get("ma7_p") is not None:
algo._entry_rows.append({
"time": algo.time,
"close": algo._last_close,
"otf_up_close": en["otf_up"],
"otf_down_close": en["otf_down"],
"ma7_c_15m": en["ma7_c"],
"ma7_p_15m": en["ma7_p"],
"ma_slope_up_15m": en["ma_slope_up"],
"ma_slope_down_15m": en["ma_slope_down"],
"ma7_c_30m": en.get("ma7_c_30m_confirm"),
"ma7_p_30m": en.get("ma7_p_30m_confirm"),
"ma_slope_up_30m": (en.get("ma7_c_30m_confirm", 0) > en.get("ma7_p_30m_confirm", 0))
if en.get("ma7_c_30m_confirm") is not None
and en.get("ma7_p_30m_confirm") is not None
else None,
"ma_slope_down_30m": (en.get("ma7_c_30m_confirm", 0) < en.get("ma7_p_30m_confirm", 0))
if en.get("ma7_c_30m_confirm") is not None
and en.get("ma7_p_30m_confirm") is not None
else None,
"confirm_long": en["confirm_long"],
"confirm_short": en["confirm_short"],
"chop_ok": en["chop_ok"],
"atr14": en["atr14"],
"stop_loss": en["stop_loss"],
"take_profit": en["take_profit"],
"entry_long": en["entry_long"],
"entry_short": en["entry_short"],
"entry_signal": en.get("entry_signal"),
"entry_confirm": en.get("entry_confirm"),
"already_long": already_long,
"already_short": already_short,
})
def handle_30m_bar(algo, bar):
# Skip missing Observix bars
bar_key = bar.end_time.strftime("%Y-%m-%d %H:%M:%S")
if bar_key in algo._missing_bars:
return
# Delegate to signal engine for any 30m state maintenance
side = algo._position_side if algo._position_side is not None else "LONG"
algo.signal_engine.update_30m_bar(bar, side, bar.end_time)
# region imports
from AlgorithmImports import *
# endregion
# indicators/chop_filter.py
class ChopFilter:
"""
Chop filter: ATR14[t] / ATR14[t-1] > threshold (ATR is expanding).
Phase 1: threshold loosened to 0.98 (from 1.0) to catch more breakout entries.
Requires ATR14 to be ready (at least 14 bars) plus one prior ATR value.
"""
def __init__(self, threshold: float = 0.98):
self._prev_atr = None
self._threshold = threshold
self.value = None
def update(self, atr14_value: float) -> None:
if atr14_value is None:
return
if self._prev_atr is not None:
if self._prev_atr == 0.0:
self.value = False
else:
self.value = (atr14_value / self._prev_atr) > self._threshold
self._prev_atr = atr14_value
def reset(self) -> None:
self._prev_atr = None
self.value = None
def is_ready(self) -> bool:
return self.value is not None# region imports
from AlgorithmImports import *
# endregion
# indicators/entry_signal.py
from collections import deque
class EntrySignal:
"""
Bar-by-bar entry signal with signal_mode support (MA7 or ZERO_LAG).
"""
def __init__(self, chop_filter: bool, confirmation_filter: bool,
signal_mode: str = "MA7", confirm_signal_mode: str = None, confirm_min_diff_pct_30m: float = 0.0):
self.chop_filter = chop_filter
self.confirmation_filter = confirmation_filter
self.signal_mode = signal_mode
# confirm_signal_mode: which signal type drives the 30m confirmation.
# Defaults to signal_mode if not specified (backward compatible).
self.confirm_signal_mode = confirm_signal_mode if confirm_signal_mode is not None else signal_mode
self.confirm_min_diff_pct_30m = float(confirm_min_diff_pct_30m)
self._prev_close = None
self._closes_15m = deque(maxlen=7)
self._prev_ma7_c = None
self.last = {"timestamp": None, "close": None, "otf_up": None, "otf_down": None,
"ma7_c": None, "ma7_p": None, "ma_slope_up": None, "ma_slope_down": None,
"chop_ok": None, "confirm_long": None, "confirm_short": None,
"entry_long": None, "entry_short": None, "atr14": None,
"stop_loss": None, "take_profit": None}
def update(self, bar, atr14_value, chop_ok, ma7_c_30m, ma7_p_30m,
sl_multiplier_long, sl_multiplier_short, tp_multiplier_long, tp_multiplier_short,
timestamp=None, sig_fast=None, sig_slow=None, sig_flat=None,
sig_diff_pct_30m=None):
close = float(bar.Close)
ts = timestamp or bar.EndTime
otf_up = (close >= self._prev_close) if self._prev_close is not None else False
otf_down = (close <= self._prev_close) if self._prev_close is not None else False
self._prev_close = close
if self.signal_mode == "ZERO_LAG":
if sig_fast is not None and sig_slow is not None:
ma7_c = sig_fast
ma7_p = sig_slow
ma_slope_up = sig_fast > sig_slow
ma_slope_down = sig_fast < sig_slow
whipsaw_flat = bool(sig_flat) if sig_flat is not None else False
else:
ma7_c = ma7_p = ma_slope_up = ma_slope_down = None
whipsaw_flat = False
else:
whipsaw_flat = False
self._closes_15m.append(close)
ma7_c = sum(self._closes_15m) / 7.0 if len(self._closes_15m) == 7 else None
if ma7_c is not None and self._prev_ma7_c is not None:
ma7_p = self._prev_ma7_c
ma_slope_up = ma7_c > ma7_p
ma_slope_down = ma7_c < ma7_p
else:
ma7_p = ma_slope_up = ma_slope_down = None
if ma7_c is not None:
self._prev_ma7_c = ma7_c
if ma7_p is None or atr14_value is None:
self._update_last(ts, close, otf_up, otf_down, ma7_c, ma7_p,
ma_slope_up, ma_slope_down, None, None, None,
False, False, atr14_value, None, None)
return False, False, None, None, None, None
chop_passes = bool(chop_ok) if (self.chop_filter and chop_ok is not None) else (not self.chop_filter)
if whipsaw_flat:
chop_passes = False
if self.confirmation_filter:
if self.confirm_signal_mode == "ZERO_LAG" and sig_diff_pct_30m is not None:
# ZL confirmation: normalized EC/EMA separation must clear an active strength threshold
confirm_long = ma_slope_up and (sig_diff_pct_30m > self.confirm_min_diff_pct_30m)
confirm_short = ma_slope_down and (sig_diff_pct_30m < -self.confirm_min_diff_pct_30m)
elif ma7_c_30m is not None and ma7_p_30m is not None:
# MA7: 30m MA7 slopes (Phase 1 behavior)
confirm_long = ma_slope_up and (ma7_c_30m > ma7_p_30m)
confirm_short = ma_slope_down and (ma7_c_30m < ma7_p_30m)
else:
confirm_long = confirm_short = False
else:
confirm_long = confirm_short = True
entry_long = bool(otf_up and ma_slope_up and chop_passes and confirm_long)
entry_short = bool(otf_down and ma_slope_down and chop_passes and confirm_short)
sl_long = close - sl_multiplier_long * atr14_value
tp_long = close + tp_multiplier_long * atr14_value
sl_short = close + sl_multiplier_short * atr14_value
tp_short = close - tp_multiplier_short * atr14_value
self._update_last(ts, close, otf_up, otf_down, ma7_c, ma7_p,
ma_slope_up, ma_slope_down, chop_passes, confirm_long, confirm_short,
entry_long, entry_short, atr14_value, sl_long, tp_long)
return entry_long, entry_short, sl_long, tp_long, sl_short, tp_short
def _update_last(self, ts, close, otf_up, otf_down, ma7_c, ma7_p,
ma_slope_up, ma_slope_down, chop_ok, confirm_long, confirm_short,
entry_long, entry_short, atr14, stop_loss, take_profit):
self.last.update({"timestamp": ts, "close": close, "otf_up": otf_up, "otf_down": otf_down,
"ma7_c": round(ma7_c, 6) if ma7_c is not None else None,
"ma7_p": round(ma7_p, 6) if ma7_p is not None else None,
"ma_slope_up": ma_slope_up, "ma_slope_down": ma_slope_down,
"chop_ok": chop_ok, "confirm_long": confirm_long, "confirm_short": confirm_short,
"entry_long": entry_long, "entry_short": entry_short,
"atr14": round(atr14, 6) if atr14 is not None else None,
"stop_loss": round(stop_loss, 4) if stop_loss is not None else None,
"take_profit": round(take_profit, 4) if take_profit is not None else None})# indicators/exit_signal.py
from AlgorithmImports import *
from collections import deque
class ExitSignal:
"""
Bar-by-bar exit signal with signal_mode support (MA7 or ZERO_LAG).
LONG exit = ~otf_up AND (flat OR ~direction)
SHORT exit = ~otf_down AND (flat OR ~direction)
"""
COLUMNS = ["timestamp", "otf_up", "otf_down", "ma7_c", "ma7_p", "flat", "direction", "exit"]
def __init__(self, threshold, warm_up_bars=14, label="", signal_mode="MA7"):
self.threshold = threshold
self.warm_up_bars = max(warm_up_bars, 8)
self.label = f"_{label}" if label else ""
self.signal_mode = signal_mode
self._closes = deque(maxlen=7)
self._prev_high = None
self._prev_low = None
self._bar_count = 0
self._prev_ma7_c = None
self.last = {k: None for k in self.COLUMNS}
def update(self, bar, position_side, timestamp=None, sig_fast=None, sig_slow=None, sig_flat=None):
high = float(bar.High)
low = float(bar.Low)
close = float(bar.Close)
ts = timestamp or bar.EndTime
self._bar_count += 1
if self._prev_high is not None:
otf_up = (low >= self._prev_low) and (position_side == "LONG")
otf_down = (high <= self._prev_high) and (position_side == "SHORT")
else:
otf_up = True
otf_down = True
self._prev_high = high
self._prev_low = low
if self.signal_mode == "ZERO_LAG":
if sig_fast is None or sig_slow is None:
return False
if self._bar_count < self.warm_up_bars:
return False
ma7_c = sig_fast
ma7_p = sig_slow
flat = bool(sig_flat) if sig_flat is not None else False
direction = (ma7_c > ma7_p) if position_side == "LONG" else (ma7_c <= ma7_p)
else:
self._closes.append(close)
if self._bar_count < self.warm_up_bars or len(self._closes) < 7:
if len(self._closes) == 7:
self._prev_ma7_c = sum(self._closes) / 7.0
return False
ma7_c = sum(self._closes) / 7.0
ma7_p = self._prev_ma7_c
self._prev_ma7_c = ma7_c
if ma7_p is None or ma7_p == 0.0:
return False
flat = abs(1.0 - ma7_c / ma7_p) < self.threshold
direction = (ma7_c > ma7_p) if position_side == "LONG" else (ma7_c <= ma7_p)
if position_side == "LONG":
exit_signal = (not otf_up) and (flat or not direction)
else:
exit_signal = (not otf_down) and (flat or not direction)
self._update_last(ts, position_side, otf_up, otf_down, ma7_c, ma7_p, flat, direction, exit_signal)
return exit_signal
def reset(self):
self._closes.clear()
self._prev_high = None
self._prev_low = None
self._bar_count = 0
self._prev_ma7_c = None
for k in self.last:
self.last[k] = None
def is_ready(self):
return self._bar_count >= self.warm_up_bars
def _update_last(self, ts, side, otf_up, otf_down, ma7_c, ma7_p, flat, direction, exit_signal):
self.last.update({"timestamp": ts, "otf_up": otf_up, "otf_down": otf_down,
"ma7_c": round(ma7_c, 6), "ma7_p": round(ma7_p, 6),
"flat": flat, "direction": direction, "exit": exit_signal})def configured_fee_per_side(config_or_fee, default=-2.25) -> float:
if isinstance(config_or_fee, dict):
return float(config_or_fee.get("trade", {}).get("fees", default))
return float(config_or_fee)
def round_trip_fee_adjustment(config_or_fee, contract_size: int = 1, default=-2.25) -> float:
return 2 * configured_fee_per_side(config_or_fee, default) * contract_size
def round_trip_fee_cost(config_or_fee, contract_size: int = 1, default=-2.25) -> float:
return abs(round_trip_fee_adjustment(config_or_fee, contract_size, default))
def apply_round_trip_fee_adjustment(gross_pnl: float, config_or_fee,
contract_size: int = 1,
default=-2.25) -> float:
return gross_pnl + round_trip_fee_adjustment(config_or_fee, contract_size, default)
def net_pnl_after_fees(gross_pnl: float, config_or_fee,
contract_size: int = 1,
default=-2.25) -> float:
return gross_pnl - round_trip_fee_cost(config_or_fee, contract_size, default)
# indicators/futures_execution.py
# Single source of truth for mapped-future execution invariants.
from AlgorithmImports import *
from indicators.session_slippage import SessionSlippageModel
class MappedFutureExecution:
"""Owns mapped subscriptions and the identity of the contract actually held."""
def __init__(self, algorithm, continuous_future,
rth_slippage_points=0.25, extended_slippage_points=1.5):
self._algorithm = algorithm
self._future = continuous_future
self._last_mapped = None
self._held_symbol = None
self._slippage_model = SessionSlippageModel(
algorithm, rth_slippage_points, extended_slippage_points
)
algorithm.add_security_initializer(self.initialize_security)
self.initialize_security(continuous_future)
def initialize_security(self, security):
"""Apply execution realism to the canonical and every mapped security."""
security.set_slippage_model(self._slippage_model)
def on_securities_changed(self, changes):
for security in changes.added_securities:
self.initialize_security(security)
def ensure_mapped_subscription(self):
"""Subscribe once when the front-month mapping changes."""
mapped = self._future.mapped
if mapped is None or mapped == self._last_mapped:
return
self._algorithm.add_future_contract(
mapped, Resolution.MINUTE, extended_market_hours=True
)
self._last_mapped = mapped
self._algorithm.log(
f"roll fix: subscribed {mapped.value} at {self._algorithm.time}"
)
def can_submit_entry(self):
"""Reject entries that cannot be submitted cleanly."""
mapped = self._future.mapped
if mapped is None:
return False
if self._algorithm.transactions.get_open_orders(mapped):
return False
return self._algorithm.is_market_open(mapped)
def record_entry(self):
"""Freeze the mapped symbol used by the entry order."""
self._held_symbol = self._future.mapped
return self._held_symbol
def exit_symbol(self):
"""Always exit the contract entered, even if mapping rolled meanwhile."""
return self._held_symbol or self._future.mapped
def clear_position(self):
self._held_symbol = None
# region imports
from AlgorithmImports import *
# endregion
# indicators/hmm_regime.py
# HMM Regime Filter -- reads pre-computed regime classifications from ObjectStore.
# Falls back to embedded BAD dates when ObjectStore key is missing.
# GPU generates the regime table with hmmlearn, QC uses it as a lookup.
import json
# 62 BAD dates from GPU HMM regime generator (2015-2026, updated 2026-04-05)
_EMBEDDED_BAD_DATES = set(['2019-08-23', '2020-02-25', '2020-02-27', '2020-03-08', '2020-03-11', '2020-03-12', '2020-03-15', '2020-03-16', '2020-03-17', '2020-03-18', '2020-03-20', '2020-03-31', '2020-04-30', '2020-06-11', '2020-07-23', '2020-09-03', '2020-09-08', '2020-09-16', '2020-11-09', '2021-01-27', '2021-02-25', '2021-03-03', '2021-05-10', '2021-12-16', '2022-01-05', '2022-01-20', '2022-01-24', '2022-01-26', '2022-02-21', '2022-02-23', '2022-03-07', '2022-04-21', '2022-04-26', '2022-04-29', '2022-05-05', '2022-05-11', '2022-05-18', '2022-06-10', '2022-06-16', '2022-06-28', '2022-08-26', '2022-09-13', '2022-09-21', '2022-09-29', '2022-10-07', '2022-10-14', '2022-10-27', '2022-11-02', '2022-12-15', '2022-12-22', '2023-08-24', '2024-08-01', '2024-09-03', '2024-12-18', '2025-03-10', '2025-04-02', '2025-04-04', '2025-04-06', '2025-04-08', '2025-04-10', '2025-10-10', '2025-11-20'])
class HMMRegimeFilter:
"""
Classifies each trading day as GOOD or BAD based on pre-computed HMM states.
Priority:
1. ObjectStore regime table (full date-indexed JSON)
2. Embedded BAD dates (fallback for paper trading / backtesting)
"""
def __init__(self, algorithm, object_store_key='hmm_regime_table', default_regime='GOOD'):
self._algo = algorithm
self._default = default_regime
self._table = {}
self._enabled = False
# Try ObjectStore first
if algorithm.object_store.contains_key(object_store_key):
raw = algorithm.object_store.read(object_store_key)
try:
self._table = json.loads(raw)
self._enabled = True
good = sum(1 for v in self._table.values() if v == 'GOOD')
bad = sum(1 for v in self._table.values() if v == 'BAD')
algorithm.log(f'HMM regime: ObjectStore loaded ({good} GOOD, {bad} BAD)')
except Exception as e:
algorithm.log(f'HMM regime: ObjectStore parse error: {e}')
if not self._enabled:
# Fall back to embedded BAD dates
self._enabled = True
algorithm.log(f'HMM regime: using embedded table ({len(_EMBEDDED_BAD_DATES)} BAD dates)')
@property
def enabled(self):
return self._enabled
def is_good_regime(self, date):
if not self._enabled:
return True
if hasattr(date, 'date'):
date = date.date()
date_str = str(date)
# ObjectStore table takes priority
if self._table:
regime = self._table.get(date_str, self._default)
return regime == 'GOOD'
# Embedded BAD dates fallback
return date_str not in _EMBEDDED_BAD_DATES
def get_regime(self, date):
if not self._enabled:
return 'DISABLED'
if hasattr(date, 'date'):
date = date.date()
date_str = str(date)
if self._table:
return self._table.get(date_str, 'UNKNOWN')
return 'BAD' if date_str in _EMBEDDED_BAD_DATES else 'GOOD'
def reset(self):
pass
from AlgorithmImports import *
from datetime import time
from indicators.fee_logic import net_pnl_after_fees, round_trip_fee_cost
def open_position(algo, side: str, price: float,
stop_loss: float, take_profit: float,
timestamp, bar) -> None:
"""
v2.5: places market order, sets local state, seeds trade tracker with BAR CLOSE
(will be overwritten with actual fill price in on_order_event). Entry Telegram
is NOT sent here -- it's sent in on_order_event once the fill price is known.
"""
# Centralized execution guard: mapped, no open order, market open.
if not algo._execution.can_submit_entry():
return
if algo.portfolio.invested:
algo.log("_open_position called while already invested -- skipped.")
return
contract_size = algo.config["trade"]["contract_size"]
# Set local state first so on_order_event has context on the fill
algo._position_side = side
algo._execution.record_entry()
algo._stop_loss = stop_loss
algo._take_profit = take_profit
# Seed trade tracker with bar close -- entry_price will be overwritten with
# actual fill price in on_order_event. MFE/MAE tracking uses current bar h/l.
algo.trade_tracker.open_trade(side, price, stop_loss, take_profit, timestamp)
algo.trade_tracker.seed_entry_bar(bar)
# Increment tag sequence (reset happens at top of _on_15m_bar on date change)
today = timestamp.date() if hasattr(timestamp, 'date') else timestamp
algo._daily_trade_seq += 1
tag = f"OTFMA_NQ_{today.strftime('%Y%m%d')}_{algo._daily_trade_seq:03d}"
algo._current_trade_tag = tag
regime_label = "GOOD"
if algo.hmm_regime is not None:
regime_label = "GOOD" if algo.hmm_regime.is_good_regime(timestamp) else "BAD"
ppp = algo.config["trade"]["price_per_point"]
t = timestamp.time() if hasattr(timestamp, 'time') else timestamp
phase = "RTH" if time(9, 30) <= t <= time(16, 0) else "EXT"
# Derive contract month from mapped symbol (v2.5 fix)
# str(algo.nq.mapped) returns QC's internal ID like "NQ Z3DIALNO785D" (last 3 = "85D" wrong).
# algo.nq.mapped.value returns ticker like "NQ18M26" (last 3 = "M26" correct).
try:
sym_obj = algo.nq.mapped
sym_str = getattr(sym_obj, 'value', None) or getattr(sym_obj, 'Value', None) or str(sym_obj)
month_code = sym_str[-3:] if len(sym_str) >= 3 else "???"
except Exception:
month_code = "???"
# Store pending entry context. on_order_event will attach fill_price + send Telegram.
algo._pending_entry = {
"signal_type": "ENTRY",
"direction": side,
"trade_tag": tag,
"contract_month": month_code,
"contract_size": contract_size,
"price_per_point": ppp,
"session_phase": phase,
"regime": regime_label,
"bar_close_price": price, # fallback if on_order_event never fires
}
# Place the market order -- fill event will drive Telegram send
if side == "LONG":
algo.market_order(algo.nq.mapped, contract_size)
else:
algo.market_order(algo.nq.mapped, -contract_size)
def close_position(algo, exit_price: float,
exit_dt, reason: str) -> None:
"""
v2.5: places reversing market order and stores pending exit context. The actual
trade_tracker.close_trade() call, PnL computation, and Telegram send all happen
in on_order_event() once the real exit fill price is known.
"""
if not algo.portfolio.invested:
return
contract_size = algo.config["trade"]["contract_size"]
tag = algo._current_trade_tag or "OTFMA_NQ_UNKNOWN"
# Store context for on_order_event to complete the close
algo._pending_exit = {
"direction": algo._position_side,
"reason": reason,
"exit_dt": exit_dt,
"bar_close_exit_price": exit_price, # fallback
"tag": tag,
"contract_size": contract_size,
}
# Exit the exact security opened; mapping may have rolled since entry.
close_symbol = algo._execution.exit_symbol()
if algo._position_side == "LONG":
algo.market_order(close_symbol, -contract_size)
elif algo._position_side == "SHORT":
algo.market_order(close_symbol, contract_size)
else:
algo.log("WARNING: _close_position with unknown side, falling back to liquidate")
algo.liquidate(close_symbol)
def handle_order_event(algo, order_event: OrderEvent):
"""
Fill confirmation + broker-state reconciliation on every order event.
Live-mode only -- backtests use synthetic fills with no discrepancies.
Prevents state drift between algo._position_side and actual broker holdings.
"""
if not algo.live_mode:
return
status = order_event.status
if status == OrderStatus.FILLED:
try:
qty = abs(order_event.fill_quantity)
fill_price = float(order_event.fill_price)
algo.log(f"ORDER FILLED: {order_event.direction} {qty} @ {fill_price}")
except Exception as e:
algo.log(f"ORDER FILLED: (log error: {e})")
fill_price = None
# v2.5: Handle pending ENTRY fill (send Telegram with real fill price)
if algo._pending_entry is not None and fill_price is not None:
p = algo._pending_entry
# Overwrite trade tracker's entry_price with the actual fill price.
# This makes MFE/MAE and any subsequent PnL use the real entry, not bar close.
if algo.trade_tracker.is_open:
algo.trade_tracker.entry_price = fill_price
# Also reset best/worst to the fill price as starting reference
algo.trade_tracker._best_price = fill_price
algo.trade_tracker._worst_price = fill_price
notional = fill_price * p["price_per_point"] * p["contract_size"]
equity = algo.portfolio.total_portfolio_value
t_now = algo.time.strftime("%H:%M ET")
# algo.telegram.notify_telegram({
# "signal_type": "ENTRY",
# "direction": p["direction"],
# "price": fill_price,
# "trade_tag": p["trade_tag"],
# "contract_month": p["contract_month"],
# "contract_size": p["contract_size"],
# "price_per_point": p["price_per_point"],
# "notional": notional,
# "session_phase": p["session_phase"],
# "regime": p["regime"],
# "equity": equity,
# })
algo._pending_entry = None
# v2.5: Handle pending EXIT fill (close tracker w/ fill price, proper PnL math, Telegram)
elif algo._pending_exit is not None and fill_price is not None:
p = algo._pending_exit
contract_size = p["contract_size"]
ppp = algo.config["trade"]["price_per_point"]
# Close the trade tracker with actual fill price (was seeded w/ bar close)
record = algo.trade_tracker.close_trade(fill_price, p["exit_dt"], p["reason"])
algo._trade_log.append(record)
entry_fill_price = record["Entry Price"] # overwritten to fill price on entry
# v2.5 CORRECT gross/net math (was swapped + double-fees in v2.4)
# Compute from real fill prices, no dependence on trade_tracker's mixed PnL.
pts = (fill_price - entry_fill_price) if p["direction"] == "LONG" else (entry_fill_price - fill_price)
gross_pnl = pts * ppp * contract_size # before fees
fees_total = round_trip_fee_cost(algo.config, contract_size) # round-trip, positive
net_pnl = net_pnl_after_fees(gross_pnl, algo.config, contract_size)
mfe = record["MFE"]
mae = record["MAE"]
give_back_raw = mfe - gross_pnl # always >= 0 (peak surrendered)
algo._cumulative_net_pnl += net_pnl
algo._session_pnl += net_pnl
algo._session_trades += 1
bars = record["Bars"]
hold_min = bars * 15
hold_h, hold_m = hold_min // 60, hold_min % 60
hold_str = f"{hold_h}h {hold_m}m" if hold_h > 0 else f"{hold_m}m"
today = p["exit_dt"].date() if hasattr(p["exit_dt"], "date") else p["exit_dt"]
algo._all_trades.append({
"date": today,
"net_pnl": net_pnl,
"gross_pnl": gross_pnl,
"tag": p["tag"],
})
# algo.telegram.notify_telegram({
# "signal_type": "EXIT",
# "direction": p["direction"],
# "entry_price": entry_fill_price,
# "exit_price": fill_price,
# "points": round(pts, 2),
# "hold_time": hold_str,
# "exit_reason": p["reason"],
# "trade_tag": p["tag"],
# "gross_pnl": round(gross_pnl, 2), # TRUE gross
# "fees": round(fees_total, 2), # always positive for display
# "net_pnl": round(net_pnl, 2), # TRUE net
# "cumulative_net_pnl": round(algo._cumulative_net_pnl, 2),
# "give_back": round(give_back_raw, 2), # positive internally; format as negative
# "mae": mae,
# "mfe": mfe,
# })
algo._position_side = None
algo._stop_loss = None
algo._take_profit = None
algo._pending_exit = None
# Existing reconcile logic (safety net if state somehow got out of sync)
try:
holdings = algo.portfolio[algo.nq.mapped]
broker_qty = holdings.quantity if holdings else 0
except Exception as e:
algo.log(f"on_order_event: broker lookup error: {e}")
return
if broker_qty == 0 and algo._position_side is not None:
algo.log(f"RECONCILE: broker flat but state={algo._position_side}. Resetting.")
algo._position_side = None
algo._stop_loss = None
algo._take_profit = None
if algo.trade_tracker.is_open:
algo.trade_tracker.reset()
elif broker_qty != 0 and algo._position_side is None:
side = "LONG" if broker_qty > 0 else "SHORT"
algo.log(f"RECONCILE: broker has {broker_qty} but state=None. Setting {side}.")
algo._position_side = side
elif status in (OrderStatus.CANCELED, OrderStatus.INVALID):
algo.log(f"ORDER {status}: {order_event}")
# If order failed, clear pending telegram state to avoid stale payloads
algo._pending_entry = None
algo._pending_exit = None
try:
invested = algo.portfolio[algo.nq.mapped].invested
except Exception:
invested = False
if not invested and algo._position_side is not None:
algo.log("RECONCILE: order canceled/invalid, resetting to flat")
algo._position_side = None
algo._stop_loss = None
algo._take_profit = None
if algo.trade_tracker.is_open:
algo.trade_tracker.reset()
from datetime import time
def is_rth(algo) -> bool:
"""Check if current time is during Regular Trading Hours (9:30-16:00 ET)."""
t = algo.time.time()
return time(9, 30) <= t <= time(16, 0)
def daily_reset(algo) -> None:
"""Fires at 9:30 AM ET -- reset daily counters, capture day-start equity."""
algo._daily_trade_seq = 0
algo._day_start_equity = algo.portfolio.total_portfolio_value
today = algo.time.date()
if algo._session_date != today:
algo._session_pnl = 0.0
algo._session_trades = 0
algo._session_date = today
# indicators/session_slippage.py
# Session-aware slippage model for NQ futures.
from AlgorithmImports import *
from datetime import time
class SessionSlippageModel:
"""Constant point-value slippage, larger outside RTH than inside it."""
def __init__(self, algorithm, rth_slippage_points=0.25, extended_slippage_points=1.5):
self._algorithm = algorithm
self._rth_slippage = rth_slippage_points
self._extended_slippage = extended_slippage_points
def _is_rth(self) -> bool:
t = self._algorithm.time.time()
return time(9, 30) <= t <= time(16, 0)
def get_slippage_approximation(self, asset, order) -> float:
return self._rth_slippage if self._is_rth() else self._extended_slippage
from AlgorithmImports import *
from datetime import timedelta
from indicators.atr import ATR14
from indicators.chop_filter import ChopFilter
from indicators.futures_execution import MappedFutureExecution
from indicators.signal_engine import SignalEngine
from indicators.trade_tracker import TradeTracker
def log_config(algo, backtest_start_date: str, bt_end_date: str,
bt_start_date: str, warmup_label: str,
cfg_exit: dict, cfg_entry: dict, cfg_trade: dict) -> None:
algo.log("=== CONFIG ===")
algo.log(f"Backtest start: {backtest_start_date} ({warmup_label})")
algo.log(f"Backtest end: {bt_end_date}")
algo.log(f"Trade start: {bt_start_date}")
algo.log(f"Initial cash: {cfg_trade['initial_cash']}")
algo.log(f"Contract size: {cfg_trade['contract_size']}")
algo.log(f"Price per point: {cfg_trade['price_per_point']}")
algo.log(f"Fees per operation: {cfg_trade['fees']}")
algo.log(f"Exit threshold 15m: {cfg_exit['threshold_15m']}")
algo.log(f"Exit threshold 30m: {cfg_exit['threshold_30m']}")
algo.log(f"Warm-up bars: {cfg_exit['warm_up_bars']}")
algo.log(f"Chop filter: {cfg_entry['chop_filter']}")
cfg_se = algo.config["signal_engine"]
algo.log(f"Entry signal: {cfg_se['entry_signal']} (15m primary)")
algo.log(f"Entry confirm: {'ON: ' + cfg_se['entry_confirm_signal'] + ' (30m)' if cfg_se['entry_confirm_enabled'] else 'OFF'}")
algo.log(f"Exit signal: {cfg_se['exit_signal']} (15m)")
algo.log(f"SL/TP enabled: {algo._sl_tp_enabled}")
algo.log(f"SL multiplier long: {cfg_entry['sl_multiplier_long']}")
algo.log(f"TP multiplier long: {cfg_entry['tp_multiplier_long']}")
algo.log(f"SL multiplier short: {cfg_entry['sl_multiplier_short']}")
algo.log(f"TP multiplier short: {cfg_entry['tp_multiplier_short']}")
algo.log("=== END CONFIG ===")
def setup_brokerage(algo, cfg_trade: dict) -> None:
algo.set_cash(cfg_trade["initial_cash"])
algo.set_brokerage_model(
BrokerageName.INTERACTIVE_BROKERS_BROKERAGE,
AccountType.MARGIN
)
def setup_future(algo) -> None:
algo.nq = algo.add_future(
Futures.Indices.NASDAQ_100_E_MINI,
resolution = Resolution.MINUTE,
data_normalization_mode = DataNormalizationMode.RAW,
contract_depth_offset = 0,
extended_market_hours = True
)
algo.nq.set_filter(0, 90)
# STANDARD_C3_BASE_20260727: one owner for mapped-future execution.
algo._execution = MappedFutureExecution(algo, algo.nq, 0.25, 1.5)
def setup_consolidators(algo) -> None:
# 15m: primary bar handler for all trading logic
# 30m: kept only to maintain exit_signal_30m state -- no trade actions
algo.consolidate(algo.nq.symbol, timedelta(minutes=15), algo._on_15m_bar)
algo.consolidate(algo.nq.symbol, timedelta(minutes=30), algo._on_30m_bar)
def setup_signal_engine(algo, cfg_exit: dict, cfg_entry: dict) -> None:
algo.signal_engine = SignalEngine(
cfg_engine = algo.config["signal_engine"],
cfg_signal = algo.config.get("signal", {}),
cfg_exit = cfg_exit,
cfg_entry = cfg_entry,
)
algo.log(f"SignalEngine: {algo.signal_engine.describe()}")
def setup_indicators(algo, cfg_trade: dict) -> None:
algo.atr14 = ATR14()
algo.chop_filter = ChopFilter()
algo.trade_tracker = TradeTracker(
price_per_point = cfg_trade["price_per_point"],
fees = cfg_trade["fees"]
)
def setup_runtime_state(algo, cfg_trade: dict) -> None:
# Trade state
algo._position_side = None
algo._stop_loss = None
algo._take_profit = None
algo._startup_reconciled = False
# Deploy identity
algo._deploy_id = str(algo.project_id)
# Data accumulators
algo._last_close = 0.0
if algo.live_mode:
algo.settings.seed_initial_prices = True
algo._entry_rows = []
algo._exit_rows = []
algo._trade_log = []
# Phase 1 flags
algo._long_only = True
algo._rth_only = True
# Session-level accumulators
algo._session_pnl = 0.0
algo._session_trades = 0
algo._session_date = None
# Cumulative / signal tracking
algo._cumulative_net_pnl = 0.0
algo._daily_trade_seq = 0
algo._day_start_equity = cfg_trade["initial_cash"]
algo._all_trades = []
algo._current_trade_tag = None
# Deferred Telegram state
algo._pending_entry = None
algo._pending_exit = None
def setup_hmm(algo) -> None:
# KARA-TRUNK: NO gates. HMM disabled by definition of the trunk.
# Downstream code guards on `algo.hmm_regime is not None`, so None = no gate.
algo._hmm_enabled = False
algo.hmm_regime = None
algo.log("HMM regime filter: DISABLED (KARA-TRUNK: no gates)")
def setup_schedules(algo) -> None:
algo.schedule.on(algo.date_rules.every_day(algo.nq.symbol),
algo.time_rules.at(9, 30),
algo._daily_reset)
# region imports
from AlgorithmImports import *
# endregion
# indicators/signal_engine.py
# SignalEngine: routes entry/exit signals through the correct signal mode.
# Created for the C1+C2 Component Matrix test (Road to Live, Aug 2026).
#
# 4 master switches:
# entry_confirm_enabled : True/False — 30m confirmation on/off
# entry_signal : "MA7"/"ZL" — 15m primary entry signal type
# entry_confirm_signal : "MA7"/"ZL" — 30m confirmation type (when enabled)
# exit_signal : "MA7"/"ZL" — 15m exit signal type
#
# Primary signals always use 15m bars. Confirmation always uses 30m bars.
# Delegates to existing EntrySignal / ExitSignal / ZeroLagEC (unchanged).
from collections import deque
from indicators.entry_signal import EntrySignal
from indicators.exit_signal import ExitSignal
from indicators.zero_lag_ec import ZeroLagEC
_VALID_MODES = ("MA7", "ZL")
def _to_internal(mode):
"""Map config value ("ZL") to internal indicator string ("ZERO_LAG")."""
return "ZERO_LAG" if mode == "ZL" else mode
class SignalEngine:
"""Routes entry/exit signals through MA7 or ZL-EC.
Primary entry/exit always on 15m bars.
Confirmation (when enabled) always on 30m bars.
"""
def __init__(self, cfg_engine, cfg_signal, cfg_exit, cfg_entry):
"""
Args:
cfg_engine: dict with the 4 master switches
cfg_signal: dict with ZL-EC params (zl_length_15m, zl_length_30m, etc.)
cfg_exit: dict with exit threshold and warm_up_bars
cfg_entry: dict with chop_filter, SL/TP multipliers
"""
# -- Parse & validate switches -----------------------------------------
self.entry_confirm_enabled = bool(cfg_engine.get("entry_confirm_enabled", True))
self.entry_signal_mode = cfg_engine.get("entry_signal", "MA7").upper()
self.entry_confirm_mode = cfg_engine.get("entry_confirm_signal", "MA7").upper()
self.exit_signal_mode = cfg_engine.get("exit_signal", "MA7").upper()
assert self.entry_signal_mode in _VALID_MODES, f"entry_signal must be {_VALID_MODES}"
assert self.entry_confirm_mode in _VALID_MODES, f"entry_confirm_signal must be {_VALID_MODES}"
assert self.exit_signal_mode in _VALID_MODES, f"exit_signal must be {_VALID_MODES}"
# -- Determine which ZL-EC instances are needed ------------------------
needs_zl_entry_15m = self.entry_signal_mode == "ZL"
needs_zl_exit_15m = self.exit_signal_mode == "ZL"
needs_zl_30m = self.entry_confirm_enabled and (self.entry_confirm_mode == "ZL")
needs_ma7_30m = self.entry_confirm_enabled and (self.entry_confirm_mode == "MA7")
# -- Create ZL-EC instances (only when needed) -------------------------
self.zl_entry_15m = ZeroLagEC(
length = cfg_signal.get("zl_length_15m", 12),
gain_limit = cfg_signal.get("zl_gain_limit_15m", 22),
threshold = cfg_signal.get("zl_threshold_15m", 0.01),
) if needs_zl_entry_15m else None
cfg_exit_signal = cfg_engine.get("exit_zl", {})
self.zl_exit_15m = ZeroLagEC(
length = cfg_exit_signal.get("zl_length_15m", 12),
gain_limit = cfg_exit_signal.get("zl_gain_limit_15m", 22),
threshold = cfg_exit_signal.get("zl_threshold_15m", 0.01),
) if needs_zl_exit_15m else None
self.zl_30m = ZeroLagEC(
length = cfg_signal.get("zl_length_30m", 14),
gain_limit = cfg_signal.get("zl_gain_limit_30m", 22),
threshold = cfg_signal.get("zl_threshold_30m", 0.01),
) if needs_zl_30m else None
# -- 30m MA7 confirmation buffer (only when confirm=MA7) ---------------
self._closes_30m = deque(maxlen=8) if needs_ma7_30m else None
self.ma7_c_30m = None
self.ma7_p_30m = None
# -- 30m ZL-EC diff for confirmation (only when confirm=ZL) ------------
self.zl_diff_30m = None
self.zl_diff_pct_30m = None
# -- Create 15m entry signal -------------------------------------------
self.entry = EntrySignal(
chop_filter = cfg_entry.get("chop_filter", True),
confirmation_filter = self.entry_confirm_enabled,
signal_mode = _to_internal(self.entry_signal_mode),
confirm_signal_mode = _to_internal(self.entry_confirm_mode),
confirm_min_diff_pct_30m = cfg_signal.get("confirm_min_diff_pct_30m", 0.0),
)
# -- Create 15m exit signal --------------------------------------------
self.exit = ExitSignal(
threshold = cfg_exit.get("threshold_15m", 0.0001),
warm_up_bars = cfg_exit.get("warm_up_bars", 14),
label = "15m",
signal_mode = _to_internal(self.exit_signal_mode),
)
# -- Public state for row accumulation ---------------------------------
self.last_entry = {}
self.last_exit = {}
# -------------------------------------------------------------------------
# LOGGING
# -------------------------------------------------------------------------
def describe(self):
"""Human-readable summary of switch settings for init logging."""
confirm = f"ON:{self.entry_confirm_mode}" if self.entry_confirm_enabled else "OFF"
return (
f"entry={self.entry_signal_mode} | confirm={confirm} | "
f"exit={self.exit_signal_mode} | "
f"ZL-entry-15m={'ON' if self.zl_entry_15m else 'OFF'} | "
f"ZL-exit-15m={'ON' if self.zl_exit_15m else 'OFF'} | "
f"ZL-30m={'ON' if self.zl_30m else 'OFF'}"
)
# -------------------------------------------------------------------------
# 15m BAR UPDATE (call before check_exit / update_entry)
# -------------------------------------------------------------------------
def update_indicators(self, bar):
"""Feed bar data to ZL-EC and 30m data buffers."""
close = float(bar.Close)
# Update 15m ZL-EC
if self.zl_entry_15m is not None:
self.zl_entry_15m.update(close)
if self.zl_exit_15m is not None:
self.zl_exit_15m.update(close)
# At 30m boundaries, update 30m data
if bar.end_time.minute in (0, 30):
self._update_30m_data(bar)
# -------------------------------------------------------------------------
# EXIT SIGNAL (15m only)
# -------------------------------------------------------------------------
def check_exit(self, bar, position_side, timestamp):
"""Evaluate 15m exit signal. Returns: exit_fired (bool)."""
# Build ZL kwargs for exit (only when exit uses ZL)
zl_kwargs = {}
if self.exit_signal_mode == "ZL" and self.zl_exit_15m is not None and self.zl_exit_15m.is_ready:
zl_kwargs = {
"sig_fast": self.zl_exit_15m.sig_fast,
"sig_slow": self.zl_exit_15m.sig_slow,
"sig_flat": self.zl_exit_15m.sig_flat,
}
exit_fired = self.exit.update(bar, position_side, timestamp, **zl_kwargs)
# Build last_exit for row accumulation
ex = self.exit.last
self.last_exit = {
"otf_up": ex.get("otf_up"),
"otf_down": ex.get("otf_down"),
"ma7_c": ex.get("ma7_c"),
"ma7_p": ex.get("ma7_p"),
"flat": ex.get("flat"),
"direction": ex.get("direction"),
"exit": exit_fired,
"exit_signal": self.exit_signal_mode,
}
return exit_fired
# -------------------------------------------------------------------------
# ENTRY SIGNAL (15m primary + optional 30m confirmation)
# -------------------------------------------------------------------------
def update_entry(self, bar, atr14_value, chop_ok, timestamp, cfg_entry):
"""Evaluate 15m entry with optional 30m confirmation.
Returns: (entry_long, entry_short, sl_long, tp_long, sl_short, tp_short)
"""
# Build ZL kwargs for 15m primary (only when entry uses ZL)
zl_kwargs = {}
if self.entry_signal_mode == "ZL" and self.zl_entry_15m is not None and self.zl_entry_15m.is_ready:
zl_kwargs = {
"sig_fast": self.zl_entry_15m.sig_fast,
"sig_slow": self.zl_entry_15m.sig_slow,
"sig_flat": self.zl_entry_15m.sig_flat,
}
result = self.entry.update(
bar = bar,
atr14_value = atr14_value,
chop_ok = chop_ok,
ma7_c_30m = self.ma7_c_30m,
ma7_p_30m = self.ma7_p_30m,
sl_multiplier_long = cfg_entry["sl_multiplier_long"],
sl_multiplier_short = cfg_entry["sl_multiplier_short"],
tp_multiplier_long = cfg_entry["tp_multiplier_long"],
tp_multiplier_short = cfg_entry["tp_multiplier_short"],
timestamp = timestamp,
sig_diff_pct_30m = self.zl_diff_pct_30m,
**zl_kwargs,
)
entry_long, entry_short = result[0], result[1]
# Build last_entry for row accumulation
en = self.entry.last
confirm_label = f"ON:{self.entry_confirm_mode}" if self.entry_confirm_enabled else "OFF"
self.last_entry = {
"ma7_c": en.get("ma7_c"),
"ma7_p": en.get("ma7_p"),
"otf_up": en.get("otf_up"),
"otf_down": en.get("otf_down"),
"ma_slope_up": en.get("ma_slope_up"),
"ma_slope_down": en.get("ma_slope_down"),
"chop_ok": en.get("chop_ok"),
"confirm_long": en.get("confirm_long"),
"confirm_short": en.get("confirm_short"),
"atr14": en.get("atr14"),
"stop_loss": en.get("stop_loss"),
"take_profit": en.get("take_profit"),
"entry_long": entry_long,
"entry_short": entry_short,
# 30m confirmation data
"ma7_c_30m_confirm": self.ma7_c_30m,
"ma7_p_30m_confirm": self.ma7_p_30m,
"zl_diff_30m": self.zl_diff_30m,
"zl_diff_pct_30m": self.zl_diff_pct_30m,
# Switch metadata
"entry_signal": self.entry_signal_mode,
"entry_confirm": confirm_label,
}
return result
# -------------------------------------------------------------------------
# 30m CONSOLIDATED BAR (no-op — all 30m data handled at boundaries)
# -------------------------------------------------------------------------
def update_30m_bar(self, bar, position_side, timestamp):
pass
# -------------------------------------------------------------------------
# INTERNAL: 30m data update at :00/:30 boundaries
# -------------------------------------------------------------------------
def _update_30m_data(self, bar):
"""Update 30m-sourced data (MA7 or ZL-EC) at 30m boundaries."""
close = float(bar.Close)
# 30m MA7 (only when confirmation uses MA7)
if self._closes_30m is not None:
self._closes_30m.append(close)
if len(self._closes_30m) == 8:
closes = list(self._closes_30m)
self.ma7_c_30m = sum(closes[1:8]) / 7.0
self.ma7_p_30m = sum(closes[0:7]) / 7.0
# 30m ZL-EC (only when confirmation uses ZL)
if self.zl_30m is not None:
self.zl_30m.update(close)
if self.zl_30m.is_ready:
self.zl_diff_30m = self.zl_30m.sig_diff
slow = self.zl_30m.sig_slow
self.zl_diff_pct_30m = (100.0 * self.zl_diff_30m / abs(slow)) if slow not in (None, 0) else 0.0
from AlgorithmImports import *
def load_missing_bars(algo, key: str) -> set:
if not algo.object_store.contains_key(key):
algo.log(f"No missing bars file found under key '{key}'. "
f"All bars will be processed.")
return set()
raw = algo.object_store.read(key)
timestamps = set(line.strip() for line in raw.splitlines() if line.strip())
algo.log(f"Loaded {len(timestamps)} missing bar timestamps from ObjectStore.")
return timestamps
def save_algorithm_outputs(algo) -> None:
import pandas as pd
if algo._entry_rows:
df = pd.DataFrame(algo._entry_rows)
algo.object_store.save(
"entry_indicators",
df.to_csv(index=False)
)
algo.log(f"Saved entry_indicators: {len(df)} rows.")
if algo._trade_log:
columns = ["Entry Date", "Entry Time", "Exit Date", "Exit Time",
"Type", "Entry Price", "Exit Price",
"PnL", "MFE", "MAE", "Bars", "Reason(s)"]
df_trades = pd.DataFrame(algo._trade_log, columns=columns)
algo.object_store.save(
"trade_log",
df_trades.to_csv(index=False)
)
algo.log(f"Saved trade_log: {len(df_trades)} trades.")
if algo._exit_rows:
df_exit = pd.DataFrame(algo._exit_rows)
algo.object_store.save(
"exit_indicators",
df_exit.to_csv(index=False)
)
algo.log(f"Saved exit_indicators: {len(df_exit)} rows.")
algo.log("on_end_of_algorithm: complete.")
# region imports
from AlgorithmImports import *
# endregion
import json
import urllib.request
import urllib.error
from datetime import timedelta
class TelegramNotifier:
def __init__(self, algo):
self.algo = algo
cfg_tg = self.algo.config.get("telegram", {})
self.enabled = False # KARA-TRUNK: backtest only, no alerts
self.bot_token = cfg_tg.get("bot_token", "8797663037:AAEuDWJ9KLZe-jd3U5crb-CJFkrW_wgW0z0")
self.channel_id = cfg_tg.get("channel_id", "-1003729759003")
if self.enabled:
self.algo.log(f"Telegram notifications: ENABLED (direct Bot API)")
else:
self.algo.log("Telegram notifications: DISABLED")
def send_morning_check(self):
"""Fires at 6:00 AM ET -- prove algo is alive, show position + regime."""
today = self.algo.time.date()
regime_label = "GOOD"
regime_emoji = "\u2705"
regime_note = ""
if getattr(self.algo, 'hmm_regime', None) is not None:
if not self.algo.hmm_regime.is_good_regime(today):
regime_label = "BAD"
regime_emoji = "\u274c"
regime_note = " \u2014 no entries today"
# Position state
if self.algo._position_side is not None and getattr(self.algo, 'trade_tracker', None) and self.algo.trade_tracker.is_open:
entry_p = self.algo.trade_tracker.entry_price
unrealized = (self.algo._last_close - entry_p) * self.algo.config["trade"]["price_per_point"] * self.algo.config["trade"]["contract_size"]
if self.algo._position_side == "SHORT":
unrealized = -unrealized
pos_line = f"Position: {self.algo._position_side} 1 NQ @ {entry_p:,.2f} (held overnight) \u26a0\ufe0f"
unreal_line = f"Unrealized PnL: {'+' if unrealized >= 0 else ''}{unrealized:,.2f}"
else:
pos_line = "Position: FLAT"
unreal_line = None
# Yesterday summary
yesterday = today - timedelta(days=1)
yday_trades = [t for t in self.algo._all_trades if t["date"] == yesterday]
if yday_trades:
yday_pnl = sum(t["net_pnl"] for t in yday_trades)
yday_str = f"Yesterday: {len(yday_trades)} trade{'s' if len(yday_trades) != 1 else ''}, {'+' if yday_pnl >= 0 else ''}{yday_pnl:,.2f} net"
elif self.algo._position_side is not None:
yday_str = "Yesterday: 1 trade, still open"
else:
yday_str = "Yesterday: 0 trades"
deploy_id = getattr(self.algo, '_deploy_id', 'Unknown')
equity = self.algo.portfolio.total_portfolio_value
mode = f"{'ZL-EC+HMM' if getattr(self.algo, '_signal_mode', '') == 'ZERO_LAG' and getattr(self.algo, '_hmm_enabled', False) else getattr(self.algo, '_signal_mode', 'MA7')} | LONG only | SL/TP OFF"
lines = [
f"\u2600\ufe0f [OTF-MA MORNING CHECK] {today}",
f"Status: ONLINE \u2705",
f"Project: {deploy_id}",
pos_line,
]
if unreal_line:
lines.append(unreal_line)
lines += [
f"Regime today: {regime_label} {regime_emoji}{regime_note}",
yday_str,
f"Cumulative Net PnL: {'+' if self.algo._cumulative_net_pnl >= 0 else ''}{self.algo._cumulative_net_pnl:,.2f}",
f"Equity: {equity:,.2f}",
f"Mode: {mode}",
]
self.send_telegram_message("\n".join(lines))
def send_daily_summary(self):
"""Fires at 4:00 PM ET -- daily recap with cumulative + WTD/MTD/YTD."""
today = self.algo.time.date()
today_trades = [t for t in self.algo._all_trades if t["date"] == today]
n_trades = len(today_trades)
wins = sum(1 for t in today_trades if t["net_pnl"] > 0)
losses = n_trades - wins
day_net = sum(t["net_pnl"] for t in today_trades)
best = max((t["net_pnl"] for t in today_trades), default=0)
worst = min((t["net_pnl"] for t in today_trades), default=0)
# WTD: Monday of this week
monday = today - timedelta(days=today.weekday())
wtd = sum(t["net_pnl"] for t in self.algo._all_trades if t["date"] >= monday)
# MTD: first of month
month_start = today.replace(day=1)
mtd = sum(t["net_pnl"] for t in self.algo._all_trades if t["date"] >= month_start)
# YTD: first of year
year_start = today.replace(month=1, day=1)
ytd = sum(t["net_pnl"] for t in self.algo._all_trades if t["date"] >= year_start)
equity = self.algo.portfolio.total_portfolio_value
ppp = self.algo.config["trade"]["price_per_point"]
contract_size = self.algo.config["trade"].get("contract_size", 1)
pnl_emoji = "\U0001f7e2" if day_net >= 0 else "\U0001f534"
mode = f"{'ZL-EC+HMM' if getattr(self.algo, '_signal_mode', '') == 'ZERO_LAG' and getattr(self.algo, '_hmm_enabled', False) else getattr(self.algo, '_signal_mode', 'MA7')} | LONG only | SL/TP OFF"
lines = [
f"\U0001f4cb [OTF-MA DAILY RECAP] {today}",
f"Trades closed: {n_trades} ({wins}W / {losses}L)",
f"Day Net PnL: {'+' if day_net >= 0 else ''}{day_net:,.2f} {pnl_emoji}",
f"Best: {'+' if best >= 0 else ''}{best:,.2f} | Worst: {'+' if worst >= 0 else ''}{worst:,.2f}",
]
# v2.5: Prominent OPEN POSITION block if still holding at recap time
if self.algo._position_side is not None and getattr(self.algo, 'trade_tracker', None) and self.algo.trade_tracker.is_open:
entry_p = self.algo.trade_tracker.entry_price
unrealized_pts = (self.algo._last_close - entry_p) if self.algo._position_side == "LONG" else (entry_p - self.algo._last_close)
unrealized_pnl = unrealized_pts * ppp * contract_size
up_emoji = "\U0001f7e2" if unrealized_pnl >= 0 else "\U0001f534"
bars = self.algo.trade_tracker.bars
hold_min = bars * 15
hold_h, hold_m = hold_min // 60, hold_min % 60
hold_str = f"{hold_h}h {hold_m}m" if hold_h > 0 else f"{hold_m}m"
lines.extend([
"",
f"\u26a0\ufe0f STILL HOLDING {self.algo._position_side} {contract_size} NQ @ {entry_p:,.2f}",
f" Current: {self.algo._last_close:,.2f} ({unrealized_pts:+.2f} pts \u00b7 held {hold_str})",
f" Unrealized: {'+' if unrealized_pnl >= 0 else ''}${unrealized_pnl:,.2f} {up_emoji}",
])
else:
lines.append("")
lines.append("Position: FLAT")
lines.extend([
"",
f"Cumulative Net PnL: {'+' if self.algo._cumulative_net_pnl >= 0 else ''}{self.algo._cumulative_net_pnl:,.2f}",
f"WTD: {'+' if wtd >= 0 else ''}{wtd:,.2f} | MTD: {'+' if mtd >= 0 else ''}{mtd:,.2f} | YTD: {'+' if ytd >= 0 else ''}{ytd:,.2f}",
f"Equity: {equity:,.2f}",
"",
f"Mode: {mode}",
])
self.send_telegram_message("\n".join(lines))
def notify_telegram(self, payload: dict) -> None:
"""
Format a signal payload into a Telegram message and send via Bot API.
"""
if not self.enabled:
return
sig = payload.get("signal_type", "UNKNOWN")
if sig == "ENTRY":
d = payload.get("direction", "LONG")
tag = payload.get("trade_tag", "???")
price = payload.get("price", 0)
month = payload.get("contract_month", "???")
phase = payload.get("session_phase", "RTH")
size = payload.get("contract_size", 1)
ppp = payload.get("price_per_point", 20)
notional = payload.get("notional", 0)
regime = payload.get("regime", "???")
equity = payload.get("equity", 0)
t_str = self.algo.time.strftime("%H:%M ET")
regime_emoji = "\u2705" if regime == "GOOD" else "\u274c"
msg = (
f"\U0001f7e2 [OTF-MA {d} ENTRY]\n"
f"Tag: {tag}\n"
f"Contract: NQ {month} @ {price:,.2f}\n"
f"Time: {t_str} ({phase})\n"
f"Size: {size} \u00d7 ${ppp}/pt (${notional:,.0f} notional)\n"
f"Regime: {regime} {regime_emoji} | Equity: ${equity:,.2f}"
)
elif sig == "EXIT":
d = payload.get("direction", "LONG")
tag = payload.get("trade_tag", "???")
entry_p = payload.get("entry_price", 0)
exit_p = payload.get("exit_price", 0)
pts = payload.get("points", 0)
hold = payload.get("hold_time", "?")
reason = payload.get("exit_reason", "?")
gross = payload.get("gross_pnl", 0)
fees = payload.get("fees", 0)
net = payload.get("net_pnl", 0)
cum = payload.get("cumulative_net_pnl", 0)
gb = payload.get("give_back", 0)
mae = payload.get("mae", 0)
mfe = payload.get("mfe", 0)
pnl_emoji = "\U0001f7e2" if net >= 0 else "\U0001f534"
pts_sign = "+" if pts >= 0 else ""
fees_display = -abs(fees) # always negative
gb_display = -abs(gb) # always negative (money surrendered from MFE peak)
msg = (
f"\U0001f3f4 [OTF-MA {d} EXIT]\n"
f"Tag: {tag}\n"
f"{entry_p:,.2f} \u2192 {exit_p:,.2f} ({pts_sign}{pts:,.2f} pts)\n"
f"Hold: {hold} | Reason: {reason}\n"
f"\n"
f"\U0001f4ca Trade Performance\n"
f"Gross PnL: {'+' if gross >= 0 else ''}{gross:,.2f} {pnl_emoji}\n"
f"Fees: {fees_display:,.2f}\n"
f"Net PnL: {'+' if net >= 0 else ''}{net:,.2f} {pnl_emoji}\n"
f"\n"
f"\U0001f4ca Cumulative Performance\n"
f"Cumulative Net PnL: {'+' if cum >= 0 else ''}{cum:,.2f}\n"
f"Give-back from MFE: {gb_display:,.2f}\n"
f"MAE: {'+' if mae >= 0 else ''}{mae:,.2f} | MFE: +{mfe:,.2f}"
)
else:
msg = json.dumps(payload)
self.send_telegram_message(msg)
def send_telegram_message(self, text: str, max_retries: int = 3) -> None:
return # SANDBOX: telegram disabled
"""
Send a plain-text message to Telegram via direct Bot API.
"""
if not self.enabled:
self.algo.log("Telegram: skipped (not enabled in config)")
return
try:
live_mode_val = self.algo.live_mode
except Exception:
live_mode_val = "unknown"
self.algo.log(f"Telegram: sending ({len(text)} chars, live_mode={live_mode_val})")
url = f"https://api.telegram.org/bot{self.bot_token}/sendMessage"
body = json.dumps(
{"chat_id": self.channel_id, "text": text},
ensure_ascii=False
).encode("utf-8")
for attempt in range(1, max_retries + 1):
try:
req = urllib.request.Request(
url, data=body,
headers={"Content-Type": "application/json; charset=utf-8"},
method="POST"
)
with urllib.request.urlopen(req, timeout=10) as resp:
result = json.loads(resp.read().decode("utf-8"))
if result.get("ok", False):
self.algo.log(f"Telegram: sent OK (attempt {attempt})")
return
self.algo.log(f"Telegram: API error attempt {attempt}: {result}")
except Exception as e:
self.algo.log(f"Telegram: send failed attempt {attempt}/{max_retries}: {type(e).__name__}: {e}")
self.algo.log(f"Telegram: FAILED after {max_retries} attempts")# region imports
from AlgorithmImports import *
# endregion
# indicators/trade_tracker.py
#
# Tracks open trade state and computes MFE, MAE, PnL on close.
# fees = -2.25 per operation (-4.50 round trip), read from config.
from datetime import datetime
from indicators.fee_logic import apply_round_trip_fee_adjustment
class TradeTracker:
"""
Maintains open trade state and produces trade records matching
the close_trade_3_8() output structure.
MFE/MAE are tracked bar-by-bar from the running high/low seen
since entry (no DataFrame slice lookup -- QC is stateful).
"""
def __init__(self, price_per_point: float, fees: float):
self.price_per_point = price_per_point
self.fees = fees
self.is_open = False
self.side = None
self.entry_price = None
self.entry_dt = None
self.stop_loss = None
self.take_profit = None
self.bars = 0
self._best_price = None
self._worst_price = None
def open_trade(self, side: str, price: float,
stop_loss: float, take_profit: float,
timestamp: datetime) -> None:
self.is_open = True
self.side = side
self.entry_price = price
self.entry_dt = timestamp
self.stop_loss = stop_loss
self.take_profit = take_profit
self.bars = 1
self._best_price = price
self._worst_price = price
def seed_entry_bar(self, bar) -> None:
if not self.is_open:
return
high = float(bar.High)
low = float(bar.Low)
if self.side == "LONG":
self._best_price = max(self._best_price, high)
self._worst_price = min(self._worst_price, low)
else:
self._best_price = min(self._best_price, low)
self._worst_price = max(self._worst_price, high)
def update_bar(self, bar) -> tuple:
if not self.is_open:
return False, False, None, ""
self.bars += 1
high = float(bar.High)
low = float(bar.Low)
close = float(bar.Close)
if self.side == "LONG":
self._best_price = max(self._best_price, high)
self._worst_price = min(self._worst_price, low)
else:
self._best_price = min(self._best_price, low)
self._worst_price = max(self._worst_price, high)
sl_hit = tp_hit = False
exit_price = close
reason = ""
if self.side == "LONG":
if self.take_profit is not None and close >= self.take_profit:
tp_hit = True
exit_price = close
reason = f"Take profit at {close:.2f}"
elif self.stop_loss is not None and close <= self.stop_loss:
sl_hit = True
exit_price = close
reason = f"Stop loss hit at {close:.2f}"
else:
if self.take_profit is not None and close <= self.take_profit:
tp_hit = True
exit_price = close
reason = f"Take profit at {close:.2f}"
elif self.stop_loss is not None and close >= self.stop_loss:
sl_hit = True
exit_price = close
reason = f"Stop loss hit at {close:.2f}"
return sl_hit, tp_hit, exit_price, reason
def close_trade(self, exit_price: float,
exit_dt: datetime,
reason: str) -> dict:
if not self.is_open:
return {
"Entry Date": None, "Entry Time": None,
"Exit Date": exit_dt.date(), "Exit Time": exit_dt.time(),
"Type": "UNKNOWN", "Entry Price": 0, "Exit Price": exit_price,
"PnL": 0, "MFE": 0, "MAE": 0, "Bars": 0,
"Reason(s)": f"close_trade called with no open trade: {reason}",
}
p = self.price_per_point
if self.side == "LONG":
MFE = (self._best_price - self.entry_price) * p
MAE = (self._worst_price - self.entry_price) * p
PnL = apply_round_trip_fee_adjustment(
p * (exit_price - self.entry_price), self.fees
)
else:
MFE = (self.entry_price - self._best_price) * p
MAE = (self.entry_price - self._worst_price) * p
PnL = apply_round_trip_fee_adjustment(
p * (self.entry_price - exit_price), self.fees
)
record = {
"Entry Date": self.entry_dt.date(),
"Entry Time": self.entry_dt.time(),
"Exit Date": exit_dt.date(),
"Exit Time": exit_dt.time(),
"Type": self.side,
"Entry Price": self.entry_price,
"Exit Price": exit_price,
"PnL": round(PnL, 2),
"MFE": round(MFE, 2),
"MAE": round(MAE, 2),
"Bars": self.bars,
"Reason(s)": reason,
}
self.is_open = False
self.side = None
self.entry_price = None
self.entry_dt = None
self.stop_loss = None
self.take_profit = None
self.bars = 0
self._best_price = None
self._worst_price = None
return record
def reset(self) -> None:
self.is_open = False
self.side = None
self.entry_price = None
self.entry_dt = None
self.stop_loss = None
self.take_profit = None
self.bars = 0
self._best_price = None
self._worst_price = None
# region imports
from AlgorithmImports import *
# endregion
# indicators/zero_lag_ec.py
#
# Ehlers Zero Lag Error-Corrected (EC) indicator.
# Stateful bar-by-bar version validated against OTF_MA_3_7_ZeroLag.ipynb.
# 5/5 tests passed with zero divergence on real NQ data.
class ZeroLagEC:
def __init__(self, length=12, gain_limit=22, threshold=0.01):
self._length = length
self._gain_limit = gain_limit
self._threshold = threshold
self._alpha = 2.0 / (length + 1)
self._ema_prev = None
self._ec_prev = None
self._bar_count = 0
self.sig_fast = None
self.sig_slow = None
self.sig_diff = None
self.sig_flat = None
self.least_error = None
self.least_error_pct = None
@property
def is_ready(self):
return self._bar_count >= 2
def update(self, close):
self._bar_count += 1
alpha = self._alpha
if self._ema_prev is None:
self._ema_prev = close
self._ec_prev = close
self.sig_fast = close
self.sig_slow = close
self.sig_diff = 0.0
self.sig_flat = True
self.least_error = 0.0
self.least_error_pct = 0.0
return
ema = alpha * close + (1.0 - alpha) * self._ema_prev
best_gain = 0.0
least_err = 1e10
for v in range(-self._gain_limit, self._gain_limit + 1):
g = v / 10.0
test_ec = alpha * (ema + g * (close - self._ec_prev)) + (1.0 - alpha) * self._ec_prev
err = abs(close - test_ec)
if err < least_err:
least_err = err
best_gain = g
ec = alpha * (ema + best_gain * (close - self._ec_prev)) + (1.0 - alpha) * self._ec_prev
least_error_pct = 100.0 * least_err / close if close != 0 else 0.0
self.sig_fast = ec
self.sig_slow = ema
self.sig_diff = ec - ema
self.sig_flat = least_error_pct < self._threshold
self.least_error = least_err
self.least_error_pct = least_error_pct
self._ema_prev = ema
self._ec_prev = ec
def reset(self):
self._ema_prev = None
self._ec_prev = None
self._bar_count = 0
self.sig_fast = None
self.sig_slow = None
self.sig_diff = None
self.sig_flat = None
self.least_error = None
self.least_error_pct = Nonefrom AlgorithmImports import *
from datetime import datetime
import calendar
from indicators.bar_logic import handle_15m_bar, handle_30m_bar
from indicators.order_logic import open_position, close_position, handle_order_event
from indicators.session_logic import is_rth, daily_reset
from indicators.setup_logic import (
log_config,
setup_brokerage,
setup_consolidators,
setup_future,
setup_hmm,
setup_indicators,
setup_runtime_state,
setup_schedules,
setup_signal_engine,
)
from indicators.storage import load_missing_bars, save_algorithm_outputs
# =============================================================================
# BACKTEST SETTINGS
# =============================================================================
bt_start_date = "2020-01-01"
bt_end_date = "2025-12-31"
bt_warmup_enabled = False
bt_warmup_months = 2
def _subtract_calendar_months(dt: datetime, months: int) -> datetime:
year = dt.year
month = dt.month - months
while month <= 0:
month += 12
year -= 1
day = min(dt.day, calendar.monthrange(year, month)[1])
return dt.replace(year=year, month=month, day=day)
def _resolve_backtest_dates():
trade_start_dt = datetime.strptime(bt_start_date, "%Y-%m-%d")
if bt_warmup_enabled:
backtest_start_dt = _subtract_calendar_months(trade_start_dt, bt_warmup_months)
else:
backtest_start_dt = trade_start_dt
warmup_label = f"{bt_warmup_months}-month warm-up" if bt_warmup_enabled else "warm-up disabled"
return backtest_start_dt.strftime("%Y-%m-%d"), warmup_label, trade_start_dt
def setup_backtest_range(algo):
backtest_start_date, warmup_label, trade_start_dt = _resolve_backtest_dates()
_sy, _sm, _sd = [int(x) for x in backtest_start_date.split('-')]
_ey, _em, _ed = [int(x) for x in bt_end_date.split('-')]
algo.set_start_date(_sy, _sm, _sd)
algo.set_end_date(_ey, _em, _ed)
algo._trade_start = trade_start_dt
return backtest_start_date, warmup_label
# =============================================================================
# EMBEDDED STRATEGY CONFIG
# =============================================================================
TRUNK_CONFIG = {
"exit": {"threshold_15m": 0.0001, "threshold_30m": 0.0005,
"warm_up_bars": 14, "sl_tp_enabled": False},
"entry": {"chop_filter": True,
"sl_multiplier_long": 1.5, "tp_multiplier_long": 3.0,
"sl_multiplier_short": 1.5, "tp_multiplier_short": 3.0},
"trade": {"initial_cash": 100000, "contract_size": 1,
"price_per_point": 20, "fees": -2.25},
"signal": {"zl_length_15m": 10, "zl_gain_limit_15m": 22, "zl_threshold_15m": 0.005,
"zl_length_30m": 14, "zl_gain_limit_30m": 26, "zl_threshold_30m": 0.01, "confirm_min_diff_pct_30m": 0.005},
# Signal routing switches.
"signal_engine": {
"entry_confirm_enabled": True, # 30m confirmation on/off
"entry_signal": "ZL", # 15m primary entry: "MA7" or "ZL"
"entry_confirm_signal": "ZL", # 30m confirmation: "MA7" or "ZL"
"exit_signal": "ZL", # 15m exit: "MA7" or "ZL"
"exit_zl": {"zl_length_15m": 14, "zl_gain_limit_15m": 22, "zl_threshold_15m": 0.01},
},
"hmm": {"enabled": False},
"telegram": {"enabled": False},
}
# =============================================================================
# QUANTCONNECT ALGORITHM
# =============================================================================
class NQStrategy(QCAlgorithm):
def initialize(self):
self.set_time_zone("America/New_York")
self.config = self._load_config("strategy_config")
cfg_exit = self.config["exit"]
cfg_entry = self.config["entry"]
cfg_trade = self.config["trade"]
backtest_start_date, warmup_label = setup_backtest_range(self)
self._sl_tp_enabled = cfg_exit.get("sl_tp_enabled", False)
log_config(
self, backtest_start_date, bt_end_date, bt_start_date,
warmup_label, cfg_exit, cfg_entry, cfg_trade
)
setup_brokerage(self, cfg_trade)
setup_future(self)
setup_consolidators(self)
setup_signal_engine(self, cfg_exit, cfg_entry)
setup_indicators(self, cfg_trade)
self._missing_bars = load_missing_bars(self, "missing_bars")
setup_runtime_state(self, cfg_trade)
setup_hmm(self)
setup_schedules(self)
self.log("OTF-MA Phase 2 initialized.")
# Contract/session callbacks.
def _is_rth(self) -> bool:
return is_rth(self)
def on_securities_changed(self, changes):
self._execution.on_securities_changed(changes)
# Data callbacks.
def on_data(self, data: Slice):
pass
def _on_15m_bar(self, bar):
handle_15m_bar(self, bar)
def _on_30m_bar(self, bar):
handle_30m_bar(self, bar)
# Order callbacks.
def _open_position(self, side: str, price: float,
stop_loss: float, take_profit: float,
timestamp, bar) -> None:
open_position(self, side, price, stop_loss, take_profit, timestamp, bar)
def _close_position(self, exit_price: float,
exit_dt, reason: str) -> None:
close_position(self, exit_price, exit_dt, reason)
def on_order_event(self, order_event: OrderEvent):
handle_order_event(self, order_event)
# Storage callback.
def on_end_of_algorithm(self):
save_algorithm_outputs(self)
# Config.
def _load_config(self, key: str) -> dict:
# ObjectStore config is intentionally ignored; this baseline is self-contained.
if self.object_store.contains_key(key):
self.log(f"Config: org ObjectStore key '{key}' EXISTS but is DELIBERATELY IGNORED (KARA-TRUNK is self-contained).")
else:
self.log("Config: no org ObjectStore key found (KARA-TRUNK self-contained either way).")
self.log(
f"Config: embedded TRUNK_CONFIG in force: "
f"exit {TRUNK_CONFIG['exit']['threshold_15m']}/{TRUNK_CONFIG['exit']['threshold_30m']}, "
f"chop {TRUNK_CONFIG['entry']['chop_filter']}, "
f"entry signal {TRUNK_CONFIG['signal_engine']['entry_signal']}, HMM off."
)
return TRUNK_CONFIG
# Scheduled callbacks.
def _daily_reset(self):
daily_reset(self)