feat: ATR-based SL/TP in backtest, portfolio-aware Kelly sizing, weighted correlation dampening, adaptive volatile threshold, fix eviction race

- backtest_engine.py: positions now also exit on ATR/regime-adaptive
  STOP_LOSS/TAKE_PROFIT (mirroring signal_service.py's close_stale_trades),
  not just REVERSAL/TIME_LIMIT/END_OF_DATA -- backtest/WFO now exercises
  the same exit rule live trading actually enforces.
- trade_executor.py: Kelly sizing dampens by 1/sqrt(same_direction_open+1)
  to account for correlated risk across simultaneously open positions
  (crypto altcoins move together); opposite-direction positions don't
  dampen since they net against that risk.
- signal_scoring.py: correlation dampening between the 13 vote algorithms
  now uses per-pair weighted coefficients (StochRSI~RSI high, MFI~RSI
  moderate, etc.) instead of uniform 1/sqrt(count), so near-duplicate
  signals get dampened harder than genuinely complementary ones.
- indicator_service.py: "volatile" regime threshold is now the 90th
  percentile of a symbol's own recent ATR% history instead of one fixed
  5% cutoff shared by every symbol (BTC vs. a naturally-volatile altcoin).
- trade_executor.py: fixed a phantom-read race in the eviction path where
  two concurrent signals for the same user could both pass the
  MAX_OPEN_TRADES check before either committed -- now locks the user row
  first to serialize per-user trade-opening.

209 backend tests pass (+22).

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
Le
2026-07-04 16:58:37 +07:00
parent 662586c6bc
commit cccef51eaf
10 changed files with 747 additions and 42 deletions
+74 -1
View File
@@ -33,6 +33,7 @@ from app.services.signal_scoring import (
_apply_regime_filter,
STRONG_BUY, BUY, STRONG_SELL, SELL,
)
from app.services.risk_manager import AdaptiveSLTPOptimizer
# Min candles for warmup: BB(20) + RSI(14) + some room = 30
MIN_CANDLES = 30
@@ -48,6 +49,38 @@ MIN_CANDLES = 30
DEFAULT_TAKER_FEE_PCT = 0.001
DEFAULT_SLIPPAGE_PCT = 0.0005
# ATR/regime-adaptive stop-loss and take-profit — mirrors the exact live
# logic in signal_service.py's close_stale_trades() (AdaptiveSLTPOptimizer
# + the same percentage floors), which used to run ONLY in live trading.
# Without this, backtest/walk-forward only ever exited on REVERSAL/
# TIME_LIMIT/END_OF_DATA — a materially different (and much more lenient)
# exit rule than what live trading actually enforces, so reported
# performance didn't reflect how live positions actually get cut short.
_SLTP_OPTIMIZER = AdaptiveSLTPOptimizer()
_DEFAULT_SL_PCT = 5.0 # used when ATR isn't available yet (e.g. warmup)
_DEFAULT_TP_PCT = 10.0
_MIN_SL_PCT = 3.0 # same floors close_stale_trades() applies
_MIN_TP_PCT = 6.0
def _effective_sl_tp_pct(
entry_price: float, atr: float | None, regime: str | None, direction: str,
) -> tuple[float, float]:
"""Percentage distance from entry at which SL/TP should fire, adapted
to current ATR and market regime — see _SLTP_OPTIMIZER above."""
if not atr or atr <= 0 or entry_price <= 0:
return _DEFAULT_SL_PCT, _DEFAULT_TP_PCT
result = _SLTP_OPTIMIZER.compute_sl_tp(
atr=atr, entry_price=entry_price, regime=regime or "neutral", direction=direction,
)
if direction == "LONG":
sl_pct = (entry_price - result["stop_loss"]) * 100.0 / entry_price
tp_pct = (result["take_profit"] - entry_price) * 100.0 / entry_price
else:
sl_pct = (result["stop_loss"] - entry_price) * 100.0 / entry_price
tp_pct = (entry_price - result["take_profit"]) * 100.0 / entry_price
return max(sl_pct, _MIN_SL_PCT), max(tp_pct, _MIN_TP_PCT)
def _fill_price(mark_price: float, direction: str, is_entry: bool, slippage_pct: float) -> float:
"""Simulate a realistic market-order fill price — slippage always moves
@@ -121,6 +154,10 @@ _SHORT_WINDOW = 3
# bounded tail, not the whole growing history, for the same O(n), not
# O(n^2), reason as everything else in this module.
_REGIME_WINDOW = 30
# (kk) detect_market_regime()'s adaptive "volatile" threshold is a
# percentile of a symbol's own recent ATR% history — a bounded trailing
# window (not the whole growing history) for the same O(n) reason.
_ATR_HISTORY_WINDOW = 100
# SMC (market_structure) and divergence detection are built from pivot
# points (local highs/lows), not rolling windows — the old approach
# precomputed a single current-state snapshot (bos/choch/trend/
@@ -245,6 +282,12 @@ def _precompute_indicators(candles: list[Candle], timeframe: str) -> dict:
# set every one of the 13 algorithms would produce unfiltered.
adx_full = adx(candle_dicts_full, period=14) or {"adx": [], "plus_di": [], "minus_di": []}
atr_full = atr(candle_dicts_full, period=14) or []
# (kk) ATR% history feeds detect_market_regime()'s adaptive "volatile"
# threshold — see _ATR_HISTORY_WINDOW above.
atr_pct_full = [
(a / c * 100.0) if a and c and c > 0 else None
for a, c in zip(atr_full, close_prices_full)
]
# Pre-build MTF candles ONCE per MTF config
mtf_precomputed = []
@@ -294,6 +337,7 @@ def _precompute_indicators(candles: list[Candle], timeframe: str) -> dict:
"price_pivot_lows_full": price_pivot_lows_full,
"adx_full": adx_full,
"atr_full": atr_full,
"atr_pct_full": atr_pct_full,
"mtf_precomputed": mtf_precomputed,
}
@@ -332,6 +376,7 @@ def _compute_scores_series(
price_pivot_lows_full = precomputed["price_pivot_lows_full"]
adx_full = precomputed["adx_full"]
atr_full = precomputed["atr_full"]
atr_pct_full = precomputed["atr_pct_full"]
mtf_precomputed = precomputed["mtf_precomputed"]
macd_hist_full = macd_full.get("histogram") if macd_full else None
@@ -476,11 +521,13 @@ def _compute_scores_series(
last_atr = atr_tail[-1] if atr_tail else None
atr_pct = (last_atr / close_prices_full[i] * 100.0) if last_atr and close_prices_full[i] > 0 else None
regime_start = max(0, i + 1 - _REGIME_WINDOW)
atr_history_start = max(0, i + 1 - _ATR_HISTORY_WINDOW)
market_regime = detect_market_regime(
adx_data, bb_data, atr_pct, vb_data,
prices=close_prices_full[regime_start:i + 1],
highs=[c["high"] for c in candle_dicts_full[regime_start:i + 1]],
lows=[c["low"] for c in candle_dicts_full[regime_start:i + 1]],
atr_pct_history=atr_pct_full[atr_history_start:i + 1],
)
results.append({
@@ -492,6 +539,7 @@ def _compute_scores_series(
"adjusted_score": adjusted_score,
"confidence": confidence,
"market_regime": market_regime,
"atr": last_atr,
})
return results
@@ -518,6 +566,14 @@ def _simulate_from_scores(
adverse slippage a live market order would, instead of pricing fills at
the exact candle close for free — see DEFAULT_TAKER_FEE_PCT/
DEFAULT_SLIPPAGE_PCT above.
Open positions are also checked every candle against ATR/regime-
adaptive STOP_LOSS/TAKE_PROFIT levels (`_effective_sl_tp_pct`) — the
same exit rule `signal_service.py`'s close_stale_trades() enforces in
live trading — in addition to REVERSAL/TIME_LIMIT/END_OF_DATA. Without
this, backtest/walk-forward positions could only ever be cut short by
a fresh opposing signal or the max-hold timer, which is materially
more lenient than what live trading actually does.
"""
all_signals: list[dict] = []
trades: list[dict] = []
@@ -569,8 +625,25 @@ def _simulate_from_scores(
)
if current_position and current_position["status"] == "OPEN":
direction = current_position["direction"]
entry_price = current_position["entry_price"]
sl_pct, tp_pct = _effective_sl_tp_pct(
entry_price, entry.get("atr"), entry.get("market_regime"), direction,
)
price_diff_pct = abs(latest_close - entry_price) / entry_price * 100.0 if entry_price > 0 else 0.0
is_adverse = (direction == "LONG" and latest_close < entry_price) or (direction == "SHORT" and latest_close > entry_price)
is_favorable = (direction == "LONG" and latest_close > entry_price) or (direction == "SHORT" and latest_close < entry_price)
hold = i - current_position["entry_index"]
if hold >= max_hold_candles:
if is_adverse and price_diff_pct >= sl_pct:
_close_position(current_position, latest_close, timestamp, "STOP_LOSS", fee_pct, slippage_pct)
trades.append(current_position)
current_position = None
elif is_favorable and price_diff_pct >= tp_pct:
_close_position(current_position, latest_close, timestamp, "TAKE_PROFIT", fee_pct, slippage_pct)
trades.append(current_position)
current_position = None
elif hold >= max_hold_candles:
_close_position(current_position, latest_close, timestamp, "TIME_LIMIT", fee_pct, slippage_pct)
trades.append(current_position)
current_position = None
+11 -4
View File
@@ -500,13 +500,19 @@ async def get_indicators(
# Add ADX (Average Directional Index) + Market Regime
computed["adx_data"] = adx(candle_dicts, period=14)
atr_vals = computed.get("supertrend", {}).get("trend", None)
# Compute ATR% for regime detection
# Compute ATR% for regime detection — also keep the full history
# (not just the latest value) so detect_market_regime can judge
# "volatile" against THIS symbol's own recent ATR% distribution
# instead of one fixed cutoff shared by every symbol (fix kk).
atr_pct_history: list[float | None] = []
try:
from app.services.indicator_service import atr as _calc_atr
raw_atr = _calc_atr(candle_dicts, period=14)
last_atr = raw_atr[-1] if raw_atr and len(raw_atr) > 0 else None
last_close = close_prices[-1] if close_prices else 1
atr_pct_val = (last_atr / last_close * 100.0) if last_atr and last_close > 0 else None
atr_pct_history = [
(a / c * 100.0) if a and c and c > 0 else None
for a, c in zip(raw_atr, close_prices)
]
atr_pct_val = atr_pct_history[-1] if atr_pct_history else None
except Exception:
atr_pct_val = None
@@ -522,6 +528,7 @@ async def get_indicators(
prices=close_prices,
highs=high_prices,
lows=low_prices,
atr_pct_history=atr_pct_history,
)
computed["market_regime"] = regime
+41 -1
View File
@@ -1114,6 +1114,40 @@ def adx(candles: list[dict], period: int = 14) -> dict[str, list[Optional[float]
# Market Regime Detection
# ======================================================================
# ── (kk) Adaptive volatile-regime threshold ──
# A fixed `atr_pct > 5.0` cutoff misclassifies regime for most symbols:
# BTC/ETH on 15m candles rarely exceed 1-2% ATR (so 5% would almost never
# fire, "volatile" never gets detected for majors), while a thin-liquidity
# altcoin or memecoin routinely trades 5%+ as its NORMAL state (so 5%
# would almost always fire, permanently downgrading/suppressing its
# signals). Percentile-based on the symbol's OWN recent ATR% history
# adapts automatically instead.
_MIN_ATR_HISTORY_FOR_ADAPTIVE_THRESHOLD = 20
_VOLATILE_PERCENTILE = 0.90
_MIN_VOLATILE_THRESHOLD_PCT = 1.0 # floor: don't flag "volatile" on trivial upticks for an unusually calm symbol
_MAX_VOLATILE_THRESHOLD_PCT = 15.0 # ceiling: sanity bound against a data glitch skewing the whole history
_DEFAULT_VOLATILE_THRESHOLD_PCT = 5.0 # fallback when there's no history yet (new symbol, cold cache)
def _adaptive_volatile_threshold(
atr_pct_history: list[Optional[float]] | None,
default: float = _DEFAULT_VOLATILE_THRESHOLD_PCT,
) -> float:
"""The `atr_pct` cutoff above which `detect_market_regime` calls the
market "volatile" — the 90th percentile of a symbol's own recent ATR%
history, so each symbol is judged against its own normal behavior
instead of one hardcoded number. Falls back to `default` when there
isn't enough history yet.
"""
if not atr_pct_history:
return default
valid = sorted(v for v in atr_pct_history if v is not None and v > 0)
if len(valid) < _MIN_ATR_HISTORY_FOR_ADAPTIVE_THRESHOLD:
return default
idx = min(int(len(valid) * _VOLATILE_PERCENTILE), len(valid) - 1)
return max(_MIN_VOLATILE_THRESHOLD_PCT, min(valid[idx], _MAX_VOLATILE_THRESHOLD_PCT))
def detect_market_regime(
adx_data: dict[str, list[Optional[float]]],
bb: dict[str, list[Optional[float]]],
@@ -1123,6 +1157,7 @@ def detect_market_regime(
highs: list[float] | None = None,
lows: list[float] | None = None,
lookback: int = 20,
atr_pct_history: list[Optional[float]] | None = None,
) -> str:
"""Classify the current market regime using multi-factor analysis.
@@ -1134,6 +1169,10 @@ def detect_market_regime(
- "breakout" : BB squeeze + volume spike
- "squeeze" : BB very narrow, low volatility before breakout
- "choppy" : high CHOP, low ER — completely avoid
`atr_pct_history`, if given, adapts the "volatile" cutoff to this
symbol's own recent ATR% distribution instead of one fixed number
shared by every symbol — see `_adaptive_volatile_threshold`.
"""
adx_vals = adx_data.get("adx", [])
adx_last = adx_vals[-1] if adx_vals and len(adx_vals) >= 1 else None
@@ -1183,7 +1222,8 @@ def detect_market_regime(
if vol_last is True:
vol_spike = True
is_volatile = atr_pct is not None and atr_pct > 5.0
volatile_threshold = _adaptive_volatile_threshold(atr_pct_history)
is_volatile = atr_pct is not None and atr_pct > volatile_threshold
# ── Enhanced classification ──
# Priority: squeeze → breakout → choppy → volatile → trending → sideways → neutral
+63 -24
View File
@@ -37,6 +37,56 @@ SQUEEZE_ALERT = "SQUEEZE_ALERT"
# Minimum distance from BB bounds to filter noise
MIN_BB_DISTANCE_PCT = Decimal("0.001") # 0.1%
# ---------------------------------------------------------------------------
# P1-16 (jj): Pairwise-weighted correlation dampening
# ---------------------------------------------------------------------------
# Strategies in the same group are correlated; dampen when multiple group
# members agree (same sign) to avoid overconfidence.
#
# A uniform 1/sqrt(count) dampening treats every co-active pair in a group
# as equally, fully correlated — but that's not true within a group:
# StochRSI is *literally derived from* RSI (empirically correlated >0.8),
# while MFI (volume-weighted RSI) typically correlates much less (~0.5)
# with either. Dampening them identically over-penalizes MFI's genuinely
# complementary volume signal while under-penalizing the near-duplicate
# RSI/StochRSI pair. Same story for the trend group (MACD/SuperTrend more
# alike than either is to Ichimoku) and volume group (a volume spike and
# cumulative OBV flow measure related but distinct things).
#
# These coefficients are domain-informed estimates, not yet calibrated
# against this system's own historical signal correlations (see
# theo_doi_trading-portal_v11.md) — but they are strictly more accurate
# than assuming every co-active pair in a group is equally (~1.0)
# correlated, which the old uniform formula implicitly did. Every pair
# that can actually occur within a group is listed explicitly; the
# default is only a safety net for future additions.
CORRELATION_GROUPS: list[list[str]] = [
["double_bb_rsi", "stoch_rsi", "mfi"], # Oscillator group
["macd_crossover", "supertrend", "ichimoku"], # Trend group
["volume_breakout", "obv"], # Volume group
["divergence", "smc", "fvg", "candlestick"], # Pattern group
]
_DEFAULT_PAIR_CORRELATION = 0.5
_PAIR_CORRELATION_WEIGHTS: dict[frozenset[str], float] = {
frozenset({"double_bb_rsi", "stoch_rsi"}): 0.85, # StochRSI is RSI re-normalized
frozenset({"double_bb_rsi", "mfi"}): 0.5, # MFI = volume-weighted RSI
frozenset({"stoch_rsi", "mfi"}): 0.5,
frozenset({"macd_crossover", "supertrend"}): 0.7, # both EMA/ATR trend-following, similar lag
frozenset({"macd_crossover", "ichimoku"}): 0.5,
frozenset({"supertrend", "ichimoku"}): 0.6,
frozenset({"volume_breakout", "obv"}): 0.4, # instantaneous spike vs. cumulative flow
frozenset({"divergence", "smc"}): 0.3,
frozenset({"divergence", "fvg"}): 0.2,
frozenset({"divergence", "candlestick"}): 0.2,
frozenset({"smc", "fvg"}): 0.35, # FVG is itself an ICT/SMC concept
frozenset({"smc", "candlestick"}): 0.2,
frozenset({"fvg", "candlestick"}): 0.2,
}
def _pair_correlation(a: str, b: str) -> float:
return _PAIR_CORRELATION_WEIGHTS.get(frozenset({a, b}), _DEFAULT_PAIR_CORRELATION)
def _get_bb_values(indicators: dict) -> dict[str, list[float]] | None:
"""Extract Bollinger Band values from indicators dict."""
@@ -449,32 +499,21 @@ def _compute_adjusted_score(
for s in disabled:
raw_scores[s] = 0.0
# ── P1-16: Correlation dampening ──
# Strategies in the same group are highly correlated; dampen when
# multiple group members agree (same sign) to avoid overconfidence.
CORRELATION_GROUPS: list[list[str]] = [
["double_bb_rsi", "stoch_rsi", "mfi"], # Oscillator group
["macd_crossover", "supertrend", "ichimoku"], # Trend group
["volume_breakout", "obv"], # Volume group
["divergence", "smc", "fvg", "candlestick"], # Pattern group
]
# ── P1-16 (jj): Pairwise-weighted correlation dampening ──
# See CORRELATION_GROUPS/_PAIR_CORRELATION_WEIGHTS/_pair_correlation
# near the top of this module for the rationale.
for group in CORRELATION_GROUPS:
active = [(s, raw_scores[s]) for s in group if raw_scores[s] != 0.0]
if len(active) >= 2:
signs = [1 if v > 0 else -1 for _, v in active]
pos_count = sum(1 for s in signs if s > 0)
neg_count = sum(1 for s in signs if s < 0)
# Dampen: scale each strategy's score by 1/sqrt(count)
if pos_count >= 2:
dampen = 1.0 / (pos_count ** 0.5)
for strat, val in active:
if val > 0:
raw_scores[strat] = val * dampen
if neg_count >= 2:
dampen = 1.0 / (neg_count ** 0.5)
for strat, val in active:
if val < 0:
raw_scores[strat] = val * dampen
for strat, val in active:
same_sign_partners = [
other_strat for other_strat, other_val in active
if other_strat != strat and (other_val > 0) == (val > 0)
]
if not same_sign_partners:
continue
corr_sum = sum(_pair_correlation(strat, other_strat) for other_strat in same_sign_partners)
dampen = 1.0 / (1.0 + corr_sum) ** 0.5
raw_scores[strat] = val * dampen
# ── Apply win-rate boosting ──
boosted_scores: dict[str, float] = {}
+29
View File
@@ -8,6 +8,7 @@ from __future__ import annotations
import json
import logging
import math
from collections import defaultdict
from datetime import datetime, timezone, timedelta
from decimal import Decimal
@@ -184,6 +185,21 @@ async def execute_signal_trade(
except Exception:
pass
# ── Serialize per-user trade-opening decisions (fix ll) ──
# The `with_for_update()` on all_open_trades below only locks rows
# that already exist — it can't stop a PHANTOM read: if two STRONG
# signals for different symbols arrive for the same user at nearly
# the same time, both transactions can lock the SAME existing open
# rows, but neither transaction's lock covers the other's brand-new
# INSERT (which doesn't exist yet to be locked). Both can then read
# the same `open_count`, both pass the MAX_OPEN_TRADES check, and
# both insert — overshooting the cap. Locking the user row itself
# (freshly, not the possibly-stale `first_user` fetched earlier)
# forces concurrent calls for this user through one at a time: the
# second call blocks here until the first commits, then re-reads
# open_count fresh and sees the first call's new trade.
await db.execute(select(User).where(User.id == first_user.id).with_for_update())
# ── Hybrid eviction ──
all_open_result = await db.execute(
select(HypotheticalTrade)
@@ -282,6 +298,19 @@ async def execute_signal_trade(
avg_loss=pnl_stats.get("avg_loss", 2.0),
confidence=signal_confidence,
)
# ── Portfolio-correlation dampening ──
# Kelly above is computed as if this were the only position —
# but crypto altcoins are typically highly correlated with each
# other (and with BTC), so a book full of same-direction
# positions carries much more compounded risk on a market-wide
# move than the per-trade Kelly fractions summed naively would
# suggest. Same-direction open positions here are the ones that
# actually stack that risk (opposite-direction positions net
# against it instead); dampen by 1/sqrt(n+1), the same
# correlated-vote dampening already used for the 13-algorithm
# signal system (see signal_scoring.py's CORRELATION_GROUPS).
same_direction_open = sum(1 for t in all_open_trades if t.direction == signal_direction)
kelly_pct *= 1.0 / math.sqrt(same_direction_open + 1)
if kelly_pct > 0:
trade_size = max(trade_size * Decimal(str(kelly_pct)), Decimal("1"))
except Exception: