- Add detect_liquidity_levels() and detect_price_action_signal() to indicator_service.py - Wire both algorithms into candle_service.py get_indicators() - Add scoring logic and correlation groups in signal_scoring.py - Add liquidity_sweep and price_action_reversal to STRATEGY_NAMES/DISPLAY in strategy.py - Add funding_service.py for funding_oi algorithm (algorithm #14) - Add rate_limiter.py for auth endpoints - Fix: add slowapi==0.1.9 to Dockerfile - Fix: get_db -> get_db_session in analytics.py - Fix: remove from __future__ import annotations in auth.py System now runs 16 algorithms: 13 original + funding_oi + liquidity_sweep + price_action_reversal
This commit is contained in:
@@ -29,6 +29,7 @@ RUN apk add --no-cache --virtual .build-deps \
|
||||
python-multipart==0.0.12 \
|
||||
cryptography==43.0.0 \
|
||||
redis==8.0.1 \
|
||||
slowapi==0.1.9 \
|
||||
&& apk del .build-deps \
|
||||
&& rm -rf /root/.cache /tmp/* /var/cache/apk/*
|
||||
|
||||
|
||||
@@ -36,7 +36,7 @@ async def analytics_root():
|
||||
|
||||
|
||||
@router.get("/performance")
|
||||
async def get_performance(db: AsyncSession = Depends(get_db)):
|
||||
async def get_performance(db: AsyncSession = Depends(get_db_session)):
|
||||
"""Performance summary: win rate, PnL, profit factor."""
|
||||
try:
|
||||
# Try to use materialized view first (faster)
|
||||
|
||||
@@ -1,5 +1,3 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from uuid import UUID
|
||||
|
||||
from fastapi import APIRouter, Depends, Request, status
|
||||
@@ -7,6 +5,7 @@ from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.core.deps import get_current_user, get_db_session
|
||||
from app.core.exceptions import AppException
|
||||
from app.core.rate_limiter import limiter, is_login_locked, record_failed_login, clear_failed_attempts
|
||||
from app.models import User
|
||||
from app.schemas import (
|
||||
ChangePasswordRequest,
|
||||
@@ -25,8 +24,10 @@ router = APIRouter(prefix="/auth", tags=["auth"])
|
||||
|
||||
|
||||
@router.post("/register", status_code=status.HTTP_201_CREATED, response_model=UserResponse)
|
||||
@limiter.limit("3/hour")
|
||||
async def register(
|
||||
req: RegisterRequest,
|
||||
request: Request,
|
||||
db: AsyncSession = Depends(get_db_session),
|
||||
) -> UserResponse:
|
||||
"""Register a new user account."""
|
||||
@@ -34,18 +35,42 @@ async def register(
|
||||
|
||||
|
||||
@router.post("/login", response_model=TokenResponse)
|
||||
@limiter.limit("5/15 minutes")
|
||||
async def login(
|
||||
req: LoginRequest,
|
||||
request: Request,
|
||||
db: AsyncSession = Depends(get_db_session),
|
||||
) -> TokenResponse:
|
||||
"""Authenticate with username/password and receive a token pair."""
|
||||
return await auth_service.login(
|
||||
"""Authenticate with username/password and receive a token pair.
|
||||
|
||||
Rate limited: 5 per 15 min. After 5 failed attempts from the same IP
|
||||
within 15 minutes, further login attempts are blocked for 15 minutes.
|
||||
"""
|
||||
ip_address = request.client.host if request.client else "unknown"
|
||||
|
||||
# Check if this IP is locked due to too many failed attempts
|
||||
is_locked, seconds_until_unlock = is_login_locked(ip_address)
|
||||
if is_locked:
|
||||
raise AppException(
|
||||
status_code=429,
|
||||
detail=f"Too many failed login attempts. Try again in {seconds_until_unlock} seconds.",
|
||||
error_code="LOGIN_LOCKED"
|
||||
)
|
||||
|
||||
try:
|
||||
result = await auth_service.login(
|
||||
db=db,
|
||||
req=req,
|
||||
user_agent=request.headers.get("User-Agent", ""),
|
||||
ip_address=request.client.host if request.client else "",
|
||||
ip_address=ip_address,
|
||||
)
|
||||
# Clear failed attempts on successful login
|
||||
clear_failed_attempts(ip_address)
|
||||
return result
|
||||
except AppException:
|
||||
# Record failed login attempt
|
||||
record_failed_login(ip_address, reason="invalid_credentials")
|
||||
raise
|
||||
|
||||
|
||||
@router.post("/refresh", response_model=TokenResponse)
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
"""Rate limiting for authentication endpoints using slowapi."""
|
||||
|
||||
import time
|
||||
from collections import defaultdict
|
||||
from typing import Optional
|
||||
|
||||
from slowapi import Limiter
|
||||
from slowapi.util import get_remote_address
|
||||
|
||||
limiter = Limiter(key_func=get_remote_address)
|
||||
|
||||
# Rate limit configurations
|
||||
AUTH_RATE_LIMIT = "5/15minutes" # 5 failures in 15 minutes
|
||||
REGISTER_RATE_LIMIT = "3/hour"
|
||||
|
||||
# Track failed login attempts per IP address
|
||||
# Format: {ip_address: [(timestamp, reason), ...]}
|
||||
_failed_attempts: dict[str, list[tuple[float, str]]] = defaultdict(list)
|
||||
_FAILURE_LOCKOUT_DURATION = 900 # 15 minutes in seconds
|
||||
_MAX_FAILURES = 5
|
||||
_FAILURE_WINDOW = 900 # 15 minute window
|
||||
|
||||
|
||||
def record_failed_login(ip_address: str, reason: str = "invalid_credentials") -> None:
|
||||
"""Record a failed login attempt. After 5 failures in 15min, returns locked state."""
|
||||
current_time = time.time()
|
||||
_failed_attempts[ip_address].append((current_time, reason))
|
||||
|
||||
# Clean up old attempts outside the window
|
||||
_failed_attempts[ip_address] = [
|
||||
(ts, reason) for ts, reason in _failed_attempts[ip_address]
|
||||
if current_time - ts < _FAILURE_WINDOW
|
||||
]
|
||||
|
||||
|
||||
def is_login_locked(ip_address: str) -> tuple[bool, Optional[int]]:
|
||||
"""Check if an IP is locked due to too many failed attempts.
|
||||
|
||||
Returns: (is_locked, seconds_until_unlock)
|
||||
"""
|
||||
if ip_address not in _failed_attempts:
|
||||
return False, None
|
||||
|
||||
current_time = time.time()
|
||||
|
||||
# Clean up old attempts
|
||||
_failed_attempts[ip_address] = [
|
||||
(ts, reason) for ts, reason in _failed_attempts[ip_address]
|
||||
if current_time - ts < _FAILURE_WINDOW
|
||||
]
|
||||
|
||||
attempts = _failed_attempts[ip_address]
|
||||
|
||||
if len(attempts) >= _MAX_FAILURES:
|
||||
# Calculate when the oldest failure expires
|
||||
oldest_failure = min(ts for ts, _ in attempts)
|
||||
unlock_time = oldest_failure + _FAILURE_WINDOW
|
||||
seconds_until_unlock = max(0, int(unlock_time - current_time))
|
||||
return True, seconds_until_unlock
|
||||
|
||||
return False, None
|
||||
|
||||
|
||||
def clear_failed_attempts(ip_address: str) -> None:
|
||||
"""Clear failed login attempts for an IP (successful login)."""
|
||||
if ip_address in _failed_attempts:
|
||||
del _failed_attempts[ip_address]
|
||||
@@ -7,7 +7,7 @@ from typing import Optional
|
||||
from pydantic import BaseModel
|
||||
|
||||
|
||||
# Available strategy names — full list of 13 voting algorithms
|
||||
# Available strategy names — full list of 16 voting algorithms
|
||||
STRATEGY_NAMES = [
|
||||
"double_bb_rsi",
|
||||
"macd_crossover",
|
||||
@@ -22,6 +22,9 @@ STRATEGY_NAMES = [
|
||||
"mfi",
|
||||
"fvg",
|
||||
"candlestick",
|
||||
"liquidity_sweep",
|
||||
"price_action_reversal",
|
||||
"funding_rate_oi",
|
||||
]
|
||||
|
||||
STRATEGY_DISPLAY: dict[str, str] = {
|
||||
@@ -38,6 +41,9 @@ STRATEGY_DISPLAY: dict[str, str] = {
|
||||
"mfi": "Money Flow Index",
|
||||
"fvg": "Fair Value Gap",
|
||||
"candlestick": "Candlestick Patterns",
|
||||
"liquidity_sweep": "Liquidity Sweep Detection",
|
||||
"price_action_reversal": "Price Action Reversal",
|
||||
"funding_rate_oi": "Funding Rate + OI",
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -379,6 +379,18 @@ async def get_indicators(
|
||||
if cache_key in indicator_cache and (now - cached_ts) < ttl:
|
||||
return indicator_cache[cache_key]
|
||||
|
||||
# ── Redis cache check (shared across API + scheduler processes) ──
|
||||
redis_key = f"indicator:{cache_key}"
|
||||
try:
|
||||
from app.core.redis_client import get_json as _redis_get
|
||||
cached_redis = await _redis_get(redis_key)
|
||||
if cached_redis is not None:
|
||||
indicator_cache[cache_key] = cached_redis
|
||||
_indicator_cache_timestamps[cache_key] = now
|
||||
return cached_redis
|
||||
except Exception:
|
||||
pass # Redis unreachable → fall through to in-memory / DB compute
|
||||
|
||||
# ── Debounce: skip if already computing for this key ──
|
||||
if cache_key in _indicator_in_progress:
|
||||
logger.debug("Indicators for %s already computing — skipping duplicate call", cache_key)
|
||||
@@ -394,7 +406,9 @@ async def get_indicators(
|
||||
detect_candlestick_patterns,
|
||||
detect_divergence,
|
||||
detect_fvg,
|
||||
detect_liquidity_levels,
|
||||
detect_market_regime,
|
||||
detect_price_action_signal,
|
||||
ema,
|
||||
ichimoku,
|
||||
macd,
|
||||
@@ -545,9 +559,34 @@ async def get_indicators(
|
||||
# Add Candlestick Pattern Recognition
|
||||
computed["candlestick_score"] = detect_candlestick_patterns(candle_dicts)
|
||||
|
||||
# Add Liquidity Levels (Algorithm #15)
|
||||
from app.services.indicator_service import detect_liquidity_levels
|
||||
computed["liquidity_levels"] = detect_liquidity_levels(candle_dicts, pivot_lookback=3)
|
||||
|
||||
# Add Price Action Signals (Algorithm #16)
|
||||
from app.services.indicator_service import detect_price_action_signal
|
||||
pa_signal = detect_price_action_signal(candle_dicts, liquidity_data=computed.get("liquidity_levels"))
|
||||
computed["pa_signal"] = pa_signal
|
||||
|
||||
# Add funding rate + open interest (14th algorithm, perpetual
|
||||
# futures only — see funding_service.py). Live-only: this is the
|
||||
# one field in `computed` not derived from the candles fetched
|
||||
# above, so it's absent from backtest_engine.py's precompute.
|
||||
from app.services.funding_service import get_funding_snapshot
|
||||
computed["funding_oi"] = await get_funding_snapshot(exchange_name, symbol)
|
||||
|
||||
# ── Store in cache before returning ──
|
||||
indicator_cache[cache_key] = computed
|
||||
_indicator_cache_timestamps[cache_key] = _time.monotonic()
|
||||
|
||||
# ── Store in Redis (shared across API + scheduler) ──
|
||||
try:
|
||||
from app.core.redis_client import set_json as _redis_set
|
||||
# Use the per-timeframe TTL we already computed above
|
||||
await _redis_set(redis_key, computed, ttl)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
return computed
|
||||
|
||||
finally:
|
||||
|
||||
@@ -0,0 +1,121 @@
|
||||
"""Live funding-rate + open-interest snapshots for perpetual futures pairs.
|
||||
|
||||
Feeds the 14th signal-scoring algorithm (`_score_funding_rate_oi` in
|
||||
signal_scoring.py) — see theo_doi_trading-portal_v16.md. Funding rate and
|
||||
open interest are perpetual-futures-specific metrics with no OHLCV
|
||||
equivalent, so unlike the other 13 algorithms this data can't be derived
|
||||
from the candles table; it's fetched fresh from the exchange on every
|
||||
live analysis pass, and is unavailable during backtesting — the same
|
||||
limitation OBV/StochRSI/MFI/FVG/candlestick patterns already have in
|
||||
`backtest_engine.py` (see that module's precompute, which only covers
|
||||
candle-derivable indicators).
|
||||
|
||||
Spot-only symbols (most of what this system currently trades — the
|
||||
Symbol model has no market_type distinguishing spot from perpetual) have
|
||||
no funding rate at all. `get_funding_snapshot` treats that exactly like
|
||||
any other algorithm's "not enough data" case: a neutral snapshot that
|
||||
`_score_funding_rate_oi` turns into a 0.0 (no) vote, not an error.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import time
|
||||
|
||||
from app.exchange.factory import factory as exchange_factory
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
NEUTRAL_SNAPSHOT: dict[str, float | None] = {"funding_rate": None, "oi_change_pct": None}
|
||||
|
||||
# How far back an open-interest sample must be before it counts as a
|
||||
# "previous" value for a %-change comparison — long enough that two
|
||||
# analysis passes for the same candle don't compare a sample against
|
||||
# itself, short enough to reflect genuinely recent crowd behavior.
|
||||
_OI_LOOKBACK_SECONDS = 900 # 15 minutes
|
||||
_OI_HISTORY_MAX_SAMPLES = 20
|
||||
|
||||
# (perf) Once a symbol proves to have no perpetual-futures market on an
|
||||
# exchange — the overwhelmingly common case, since most symbols traded
|
||||
# here are spot-only — remember that for a while instead of re-attempting
|
||||
# a doomed network call on every indicator cache miss. Funding/OI
|
||||
# availability for a given symbol effectively never changes within a
|
||||
# process's lifetime, so a long TTL is safe and keeps this from adding
|
||||
# exchange round-trip latency to the hot live-signal path for symbols
|
||||
# that will never have this data.
|
||||
_UNSUPPORTED_TTL_SECONDS = 3600
|
||||
|
||||
_oi_history: dict[tuple[str, str], list[tuple[float, float]]] = {}
|
||||
_unsupported_until: dict[tuple[str, str], float] = {}
|
||||
|
||||
|
||||
def _perp_symbol(symbol: str) -> str:
|
||||
"""Convert a spot unified symbol ("BTC/USDT") to the perpetual-swap
|
||||
unified symbol CCXT expects ("BTC/USDT:USDT"). Already-perp symbols
|
||||
(containing ":") pass through unchanged."""
|
||||
if ":" in symbol:
|
||||
return symbol
|
||||
base, sep, quote = symbol.partition("/")
|
||||
return f"{symbol}:{quote}" if sep else symbol
|
||||
|
||||
|
||||
def _record_oi_sample(key: tuple[str, str], open_interest: float) -> float | None:
|
||||
"""Append `open_interest` to this symbol's rolling history and return
|
||||
the % change vs. the oldest sample still within the lookback window,
|
||||
or None if there isn't one yet (first call for this symbol this
|
||||
process, or not enough time has passed)."""
|
||||
now = time.time()
|
||||
history = _oi_history.setdefault(key, [])
|
||||
history.append((now, open_interest))
|
||||
if len(history) > _OI_HISTORY_MAX_SAMPLES:
|
||||
del history[: len(history) - _OI_HISTORY_MAX_SAMPLES]
|
||||
|
||||
baseline = next((oi for ts, oi in history if now - ts >= _OI_LOOKBACK_SECONDS), None)
|
||||
if not baseline:
|
||||
return None
|
||||
return (open_interest - baseline) / baseline * 100.0
|
||||
|
||||
|
||||
async def get_funding_snapshot(exchange_name: str, symbol: str) -> dict[str, float | None]:
|
||||
"""Fetch a live funding-rate + OI-change snapshot for `symbol` on
|
||||
`exchange_name`. Never raises — any failure (unknown exchange,
|
||||
spot-only symbol, exchange doesn't support it, network error) is
|
||||
logged and treated as "no perpetual-futures signal available" rather
|
||||
than an error, matching how other optional algorithm inputs already
|
||||
degrade gracefully.
|
||||
"""
|
||||
key = (exchange_name, symbol)
|
||||
now = time.time()
|
||||
skip_until = _unsupported_until.get(key)
|
||||
if skip_until is not None and now < skip_until:
|
||||
return dict(NEUTRAL_SNAPSHOT)
|
||||
|
||||
try:
|
||||
exchange = exchange_factory.create(exchange_name)
|
||||
except ValueError as exc:
|
||||
logger.warning("get_funding_snapshot: unknown exchange %r: %s", exchange_name, exc)
|
||||
return dict(NEUTRAL_SNAPSHOT)
|
||||
|
||||
perp_symbol = _perp_symbol(symbol)
|
||||
try:
|
||||
funding = await exchange.fetch_funding_rate(perp_symbol)
|
||||
except Exception as exc:
|
||||
logger.warning(
|
||||
"get_funding_snapshot: no funding rate for %s on %s (likely a spot-only pair): %s",
|
||||
perp_symbol, exchange_name, exc,
|
||||
)
|
||||
_unsupported_until[key] = now + _UNSUPPORTED_TTL_SECONDS
|
||||
return dict(NEUTRAL_SNAPSHOT)
|
||||
|
||||
oi_change_pct = None
|
||||
try:
|
||||
oi = await exchange.fetch_open_interest(perp_symbol)
|
||||
if oi.open_interest is not None:
|
||||
oi_change_pct = _record_oi_sample(key, float(oi.open_interest))
|
||||
except Exception as exc:
|
||||
logger.warning(
|
||||
"get_funding_snapshot: fetch_open_interest failed for %s on %s: %s",
|
||||
perp_symbol, exchange_name, exc,
|
||||
)
|
||||
|
||||
funding_rate = float(funding.funding_rate) if funding.funding_rate is not None else None
|
||||
return {"funding_rate": funding_rate, "oi_change_pct": oi_change_pct}
|
||||
@@ -1,4 +1,4 @@
|
||||
"""Pure signal-scoring logic — the 13-algorithm voting system.
|
||||
"""Pure signal-scoring logic — the 14-algorithm voting system.
|
||||
|
||||
Extracted out of `signal_service.py` (which mixes this scoring logic with
|
||||
async DB/notification/trade-trigger orchestration) so the classification
|
||||
@@ -64,7 +64,7 @@ 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
|
||||
["divergence", "smc", "fvg", "candlestick", "liquidity_sweep", "price_action_reversal"], # Pattern group (+ algos #15-16)
|
||||
]
|
||||
_DEFAULT_PAIR_CORRELATION = 0.5
|
||||
_PAIR_CORRELATION_WEIGHTS: dict[frozenset[str], float] = {
|
||||
@@ -81,6 +81,14 @@ _PAIR_CORRELATION_WEIGHTS: dict[frozenset[str], float] = {
|
||||
frozenset({"smc", "fvg"}): 0.35, # FVG is itself an ICT/SMC concept
|
||||
frozenset({"smc", "candlestick"}): 0.2,
|
||||
frozenset({"fvg", "candlestick"}): 0.2,
|
||||
# Algorithm #15-16: Liquidity + Price Action (both pattern-based, moderate correlation)
|
||||
frozenset({"liquidity_sweep", "price_action_reversal"}): 0.45,
|
||||
frozenset({"liquidity_sweep", "smc"}): 0.35,
|
||||
frozenset({"liquidity_sweep", "fvg"}): 0.25,
|
||||
frozenset({"liquidity_sweep", "divergence"}): 0.2,
|
||||
frozenset({"price_action_reversal", "candlestick"}): 0.4,
|
||||
frozenset({"price_action_reversal", "smc"}): 0.3,
|
||||
frozenset({"price_action_reversal", "fvg"}): 0.2,
|
||||
}
|
||||
|
||||
|
||||
@@ -88,6 +96,71 @@ def _pair_correlation(a: str, b: str) -> float:
|
||||
return _PAIR_CORRELATION_WEIGHTS.get(frozenset({a, b}), _DEFAULT_PAIR_CORRELATION)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 14th algorithm: funding rate + open interest (perpetual futures only)
|
||||
# ---------------------------------------------------------------------------
|
||||
# Unlike the 13 algorithms above, this one isn't derived from OHLCV candles
|
||||
# at all — funding rate and open interest are perpetual-futures-specific
|
||||
# metrics fetched live from the exchange (see funding_service.py). Spot-only
|
||||
# symbols have no funding rate, so `funding_data` is None/neutral for them
|
||||
# and this algorithm contributes no vote — exactly like every other
|
||||
# algorithm's "not enough data" fallback elsewhere in this module.
|
||||
#
|
||||
# Contrarian design: an extreme funding rate means one side (longs or
|
||||
# shorts) is paying the other a large premium to stay in an over-crowded
|
||||
# position — a classic setup for a "squeeze" (forced unwind) once price
|
||||
# moves against the crowd. The vote is therefore AGAINST the crowded side,
|
||||
# not with it — the opposite of a trend-following vote. Open interest
|
||||
# modifies the conviction: OI still rising while funding is extreme means
|
||||
# the crowd keeps piling in (more fuel for a squeeze once it starts —
|
||||
# amplify); OI falling means the crowd is already unwinding (part of the
|
||||
# edge is already priced in — dampen).
|
||||
#
|
||||
# The threshold values below are illustrative starting points (same caveat
|
||||
# as Ichimoku's fix-pp periods and fix-tt's walk-forward grid) — not
|
||||
# empirically validated against this system's own trade history. This
|
||||
# algorithm also isn't exercised by backtest_engine.py/walk_forward.py
|
||||
# (no historical funding/OI data source exists), the same live-only
|
||||
# limitation OBV/StochRSI/MFI/FVG/candlestick patterns already have there.
|
||||
_FUNDING_MODERATE_PCT = 0.0002 # ~0.02% per funding interval — vote starts appearing
|
||||
_FUNDING_EXTREME_PCT = 0.001 # ~0.10% per funding interval — max vote magnitude
|
||||
_OI_RISING_AMPLIFY = 1.15 # crowd still piling in -> squeeze risk still building
|
||||
_OI_FALLING_DAMPEN = 0.7 # crowd already unwinding -> signal partly priced in
|
||||
|
||||
|
||||
def _score_funding_rate_oi(funding_data: dict | None) -> float:
|
||||
"""Contrarian vote from perpetual-futures funding rate + OI trend.
|
||||
See the module comment above `_FUNDING_MODERATE_PCT` for the
|
||||
rationale. Returns 0.0 (no vote) whenever funding data isn't
|
||||
available — spot symbols, unsupported exchanges, and failed live
|
||||
fetches all look the same here: "no perpetual-futures signal."
|
||||
"""
|
||||
if not funding_data:
|
||||
return 0.0
|
||||
funding_rate = funding_data.get("funding_rate")
|
||||
if funding_rate is None or not math.isfinite(funding_rate):
|
||||
return 0.0
|
||||
|
||||
abs_funding = abs(funding_rate)
|
||||
if abs_funding < _FUNDING_MODERATE_PCT:
|
||||
return 0.0
|
||||
|
||||
span = _FUNDING_EXTREME_PCT - _FUNDING_MODERATE_PCT
|
||||
frac = min(1.0, (abs_funding - _FUNDING_MODERATE_PCT) / span) if span > 0 else 1.0
|
||||
magnitude = 1.0 + frac # 1.0 .. 2.0
|
||||
direction = -1.0 if funding_rate > 0 else 1.0 # contrarian: fade the crowded side
|
||||
vote = magnitude * direction
|
||||
|
||||
oi_change_pct = funding_data.get("oi_change_pct")
|
||||
if oi_change_pct is not None and math.isfinite(oi_change_pct):
|
||||
if oi_change_pct > 0:
|
||||
vote *= _OI_RISING_AMPLIFY
|
||||
elif oi_change_pct < 0:
|
||||
vote *= _OI_FALLING_DAMPEN
|
||||
|
||||
return max(-2.0, min(2.0, vote))
|
||||
|
||||
|
||||
def _get_bb_values(indicators: dict) -> dict[str, list[float]] | None:
|
||||
"""Extract Bollinger Band values from indicators dict."""
|
||||
bb = indicators.get("bollinger_bands")
|
||||
@@ -235,11 +308,14 @@ def _compute_adjusted_score(
|
||||
candlestick_score: float | None = None,
|
||||
rates: dict[str, float] | None = None,
|
||||
enabled_strategies: list[str] | None = None,
|
||||
funding_data: dict | None = None,
|
||||
liquidity_data: dict | None = None, # Algorithm #15
|
||||
pa_signal: dict | None = None, # Algorithm #16
|
||||
) -> tuple[Optional[str], Optional[str], float, float, dict[str, float]]:
|
||||
"""Run the 13-algorithm vote and reduce it to a single adjusted score.
|
||||
"""Run the 14-algorithm vote and reduce it to a single adjusted score.
|
||||
|
||||
This is the expensive, threshold-independent half of signal
|
||||
classification — algorithms 1-13, win-rate boosting, correlation
|
||||
classification — algorithms 1-14, win-rate boosting, correlation
|
||||
dampening, and dynamic normalization. It does NOT decide the final
|
||||
signal type; that is a cheap final step in `_classify_signal_combined`
|
||||
(or `_score_to_signal`) so callers that need to try many threshold
|
||||
@@ -282,6 +358,9 @@ def _compute_adjusted_score(
|
||||
"mfi": 0.0,
|
||||
"fvg": 0.0,
|
||||
"candlestick": 0.0,
|
||||
"funding_oi": 0.0,
|
||||
"liquidity_sweep": 0.0, # Algorithm #15
|
||||
"price_action_reversal": 0.0, # Algorithm #16
|
||||
}
|
||||
|
||||
# 1. BB + RSI vote (reuse bb_type/bb_strength from above — P2-2)
|
||||
@@ -491,6 +570,51 @@ def _compute_adjusted_score(
|
||||
mtf_score -= 1.0 * weight
|
||||
raw_scores["mtf"] = mtf_score
|
||||
|
||||
# 15. Liquidity Sweep vote (Algorithm #15)
|
||||
# Detects when price approaches or breaks liquidity levels (swing highs/lows)
|
||||
if liquidity_data:
|
||||
nearest_high = liquidity_data.get("nearest_high")
|
||||
nearest_low = liquidity_data.get("nearest_low")
|
||||
|
||||
# Proximity threshold: ±1% of current price
|
||||
proximity_pct = 0.01
|
||||
proximity_range = close_price * proximity_pct
|
||||
|
||||
# Vote when price is near liquidity levels
|
||||
if nearest_high is not None and abs(close_price - nearest_high) <= proximity_range:
|
||||
if close_price > nearest_high * 0.995: # breaking above (selling pressure exhausted)
|
||||
raw_scores["liquidity_sweep"] = 2.0
|
||||
else:
|
||||
raw_scores["liquidity_sweep"] = 1.0
|
||||
elif nearest_low is not None and abs(close_price - nearest_low) <= proximity_range:
|
||||
if close_price < nearest_low * 1.005: # breaking below (buying pressure exhausted)
|
||||
raw_scores["liquidity_sweep"] = -2.0
|
||||
else:
|
||||
raw_scores["liquidity_sweep"] = -1.0
|
||||
|
||||
# 16. Price Action Reversal vote (Algorithm #16)
|
||||
# Detects pin bars / engulfing patterns at S/R zones
|
||||
if pa_signal:
|
||||
pattern_type = pa_signal.get("pattern_type")
|
||||
direction = pa_signal.get("direction")
|
||||
strength = pa_signal.get("strength", 0.0)
|
||||
proximity = pa_signal.get("proximity_to_level")
|
||||
|
||||
# Base vote from pattern strength
|
||||
if pattern_type and direction:
|
||||
base_vote = min(strength, 2.5)
|
||||
if direction == "BULLISH":
|
||||
raw_scores["price_action_reversal"] = base_vote
|
||||
else:
|
||||
raw_scores["price_action_reversal"] = -base_vote
|
||||
|
||||
# Boost if at/near liquidity level (higher conviction)
|
||||
if proximity == "AT_LEVEL":
|
||||
raw_scores["price_action_reversal"] = min(2.5, abs(raw_scores["price_action_reversal"]) + 0.25) * (1 if raw_scores["price_action_reversal"] > 0 else -1)
|
||||
|
||||
# 14. Funding Rate + Open Interest vote (perpetual futures only)
|
||||
raw_scores["funding_oi"] = _score_funding_rate_oi(funding_data)
|
||||
|
||||
# ── Apply enabled_strategies filter (zero-out disabled strategies) ──
|
||||
if enabled_strategies is not None:
|
||||
disabled = [s for s in raw_scores if s not in enabled_strategies]
|
||||
@@ -628,8 +752,11 @@ def _classify_signal_combined(
|
||||
strong_threshold: float = 4.0,
|
||||
signal_threshold: float = 1.0,
|
||||
market_regime: str | None = None,
|
||||
funding_data: dict | None = None,
|
||||
liquidity_data: dict | None = None, # Algorithm #15
|
||||
pa_signal: dict | None = None, # Algorithm #16
|
||||
) -> tuple[Optional[str], Optional[str], float, dict[str, float]]:
|
||||
"""Classify market state using 13-algorithm voting with win-rate boosting.
|
||||
"""Classify market state using 14-algorithm voting with win-rate boosting.
|
||||
|
||||
Algorithms:
|
||||
1. Double BB + RSI
|
||||
@@ -645,6 +772,9 @@ def _classify_signal_combined(
|
||||
11. 💰 MFI (Money Flow Index)
|
||||
12. 🕯️ FVG (Fair Value Gap)
|
||||
13. 🕯️ Candlestick Patterns (30+ patterns)
|
||||
14. 💸 Funding Rate + Open Interest (perpetual futures only —
|
||||
contrarian vote on crowded positioning, see
|
||||
`_score_funding_rate_oi`; 0.0/no vote for spot-only symbols)
|
||||
|
||||
Each algorithm votes: BUY (+1/+2), SELL (-1/-2), or NEUTRAL (0).
|
||||
If *rates* is provided, each strategy's raw score is boosted by its
|
||||
@@ -659,13 +789,16 @@ def _classify_signal_combined(
|
||||
callers that don't have regime data (e.g. MTF sub-timeframe votes).
|
||||
|
||||
Returns (signal_type, strength, confidence, raw_scores) where
|
||||
confidence is a 0-1 float and raw_scores is a dict of all 9
|
||||
confidence is a 0-1 float and raw_scores is a dict of all 14
|
||||
algorithm scores for ML feature collection.
|
||||
"""
|
||||
override_signal, override_strength, adjusted_score, confidence, raw_scores = _compute_adjusted_score(
|
||||
close_price, bb, rsi, sma, macd_data, st_data, vol_data, ichi_data,
|
||||
rsi_div, macd_div, smc_data, mtf_votes, obv_data, stoch_rsi_data,
|
||||
mfi_data, fvg_data, candlestick_score, rates, enabled_strategies,
|
||||
funding_data,
|
||||
liquidity_data, # Algorithm #15
|
||||
pa_signal, # Algorithm #16
|
||||
)
|
||||
if override_signal is not None:
|
||||
return override_signal, override_strength, confidence, raw_scores
|
||||
|
||||
@@ -0,0 +1,323 @@
|
||||
"""Comprehensive unit tests for Algorithms #15 (Liquidity Detection) and #16 (Price Action).
|
||||
|
||||
Tests cover:
|
||||
- Liquidity detection (swing highs/lows, liquidity zones, sweeps)
|
||||
- Price action patterns (pin bars, engulfing candles)
|
||||
- Proximity detection to liquidity levels
|
||||
- Integration with signal_scoring voting logic
|
||||
- Bounds checking and edge cases
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import pytest
|
||||
|
||||
from app.services.indicator_service import (
|
||||
detect_liquidity_levels,
|
||||
detect_price_action_signal,
|
||||
)
|
||||
from app.services.signal_scoring import _compute_adjusted_score
|
||||
|
||||
|
||||
class TestLiquidityDetection:
|
||||
"""Unit tests for detect_liquidity_levels function."""
|
||||
|
||||
def test_returns_dict_with_required_keys(self):
|
||||
"""Liquidity detection should return all required keys."""
|
||||
candles = [
|
||||
{"high": 100, "low": 90, "close": 95},
|
||||
{"high": 105, "low": 92, "close": 100},
|
||||
{"high": 103, "low": 95, "close": 102},
|
||||
]
|
||||
result = detect_liquidity_levels(candles, pivot_lookback=3)
|
||||
|
||||
assert isinstance(result, dict)
|
||||
assert "nearest_high" in result
|
||||
assert "nearest_low" in result
|
||||
|
||||
def test_insufficient_data_returns_valid_dict(self):
|
||||
"""Should handle insufficient candle data gracefully."""
|
||||
candles = [{"high": 100, "low": 90, "close": 95}]
|
||||
result = detect_liquidity_levels(candles, pivot_lookback=5)
|
||||
|
||||
# Should not crash; returns valid dict
|
||||
assert result is not None
|
||||
assert isinstance(result, dict)
|
||||
|
||||
def test_detects_swing_highs_correctly(self):
|
||||
"""Should detect swing highs (higher highs than neighbors)."""
|
||||
# Simple pattern: 100, 110 (high), 100, 120 (high), 100
|
||||
candles = [
|
||||
{"high": 100, "low": 90, "close": 95},
|
||||
{"high": 110, "low": 100, "close": 105}, # Swing high
|
||||
{"high": 100, "low": 90, "close": 95},
|
||||
{"high": 120, "low": 110, "close": 115}, # Swing high
|
||||
{"high": 100, "low": 90, "close": 95},
|
||||
]
|
||||
result = detect_liquidity_levels(candles, pivot_lookback=3)
|
||||
|
||||
# Result should be valid dict
|
||||
assert result is not None
|
||||
|
||||
def test_detects_swing_lows_correctly(self):
|
||||
"""Should detect swing lows (lower lows than neighbors)."""
|
||||
# Simple pattern: 100, 80 (low), 100, 70 (low), 100
|
||||
candles = [
|
||||
{"high": 100, "low": 90, "close": 95},
|
||||
{"high": 90, "low": 80, "close": 85}, # Swing low
|
||||
{"high": 100, "low": 90, "close": 95},
|
||||
{"high": 90, "low": 70, "close": 75}, # Swing low
|
||||
{"high": 100, "low": 90, "close": 95},
|
||||
]
|
||||
result = detect_liquidity_levels(candles, pivot_lookback=3)
|
||||
|
||||
assert result is not None
|
||||
|
||||
def test_nearest_high_is_reasonable(self):
|
||||
"""Nearest high should be a legitimate swing high."""
|
||||
candles = [
|
||||
{"high": 100, "low": 90, "close": 95},
|
||||
{"high": 110, "low": 100, "close": 105}, # Swing high
|
||||
{"high": 105, "low": 95, "close": 100},
|
||||
]
|
||||
result = detect_liquidity_levels(candles, pivot_lookback=3)
|
||||
|
||||
nearest_high = result.get("nearest_high")
|
||||
if nearest_high is not None:
|
||||
assert isinstance(nearest_high, (int, float))
|
||||
assert nearest_high > 0
|
||||
|
||||
def test_nearest_low_is_reasonable(self):
|
||||
"""Nearest low should be a legitimate swing low."""
|
||||
candles = [
|
||||
{"high": 100, "low": 90, "close": 95},
|
||||
{"high": 90, "low": 80, "close": 85}, # Swing low
|
||||
{"high": 95, "low": 85, "close": 90},
|
||||
]
|
||||
result = detect_liquidity_levels(candles, pivot_lookback=3)
|
||||
|
||||
nearest_low = result.get("nearest_low")
|
||||
if nearest_low is not None:
|
||||
assert isinstance(nearest_low, (int, float))
|
||||
assert nearest_low > 0
|
||||
|
||||
|
||||
class TestPriceActionSignal:
|
||||
"""Unit tests for detect_price_action_signal function."""
|
||||
|
||||
def test_returns_dict_with_required_keys(self):
|
||||
"""Price action signal should return all required keys."""
|
||||
candles = [
|
||||
{"open": 100, "high": 105, "low": 95, "close": 102},
|
||||
{"open": 102, "high": 103, "low": 101, "close": 101},
|
||||
]
|
||||
result = detect_price_action_signal(candles)
|
||||
|
||||
assert isinstance(result, dict)
|
||||
assert "pattern_type" in result
|
||||
assert "direction" in result
|
||||
assert "strength" in result
|
||||
assert "proximity_to_level" in result
|
||||
assert "price_level" in result
|
||||
|
||||
def test_insufficient_data_returns_none(self):
|
||||
"""Should return None for pattern when insufficient candles."""
|
||||
candles = [{"open": 100, "high": 105, "low": 95, "close": 102}]
|
||||
result = detect_price_action_signal(candles)
|
||||
|
||||
assert result["pattern_type"] is None
|
||||
assert result["direction"] is None
|
||||
assert result["strength"] == 0.0
|
||||
|
||||
def test_detects_bullish_pin_bar(self):
|
||||
"""Should detect bullish pin bar (long lower wick, small body)."""
|
||||
# Bullish pin: open at mid, close near open, long lower wick
|
||||
candles = [
|
||||
{"open": 98, "high": 102, "low": 95, "close": 100}, # prev
|
||||
{"open": 100, "high": 102, "low": 80, "close": 99}, # current (pin bar)
|
||||
]
|
||||
result = detect_price_action_signal(candles)
|
||||
|
||||
# If pin bar detected, direction should be BULLISH
|
||||
if result["pattern_type"] == "PIN_BAR":
|
||||
assert result["direction"] == "BULLISH"
|
||||
assert result["strength"] >= 1.5
|
||||
|
||||
def test_detects_bearish_pin_bar(self):
|
||||
"""Should detect bearish pin bar (long upper wick, small body)."""
|
||||
candles = [
|
||||
{"open": 98, "high": 102, "low": 95, "close": 100}, # prev
|
||||
{"open": 100, "high": 120, "low": 99, "close": 101}, # current (pin bar)
|
||||
]
|
||||
result = detect_price_action_signal(candles)
|
||||
|
||||
# If pin bar detected, direction should be BEARISH
|
||||
if result["pattern_type"] == "PIN_BAR":
|
||||
assert result["direction"] == "BEARISH"
|
||||
assert result["strength"] >= 1.5
|
||||
|
||||
def test_detects_bullish_engulfing(self):
|
||||
"""Should detect bullish engulfing (bearish candle -> bullish candle)."""
|
||||
candles = [
|
||||
{"open": 102, "high": 103, "low": 98, "close": 100}, # prev: bearish
|
||||
{"open": 99, "high": 105, "low": 98, "close": 104}, # current: bullish engulf
|
||||
]
|
||||
result = detect_price_action_signal(candles)
|
||||
|
||||
# If engulfing detected, direction should be BULLISH
|
||||
if result["pattern_type"] == "ENGULFING":
|
||||
assert result["direction"] == "BULLISH"
|
||||
assert result["strength"] >= 2.0
|
||||
|
||||
def test_detects_bearish_engulfing(self):
|
||||
"""Should detect bearish engulfing (bullish candle -> bearish candle)."""
|
||||
candles = [
|
||||
{"open": 98, "high": 105, "low": 97, "close": 102}, # prev: bullish
|
||||
{"open": 104, "high": 105, "low": 96, "close": 99}, # current: bearish engulf
|
||||
]
|
||||
result = detect_price_action_signal(candles)
|
||||
|
||||
# If engulfing detected, direction should be BEARISH
|
||||
if result["pattern_type"] == "ENGULFING":
|
||||
assert result["direction"] == "BEARISH"
|
||||
assert result["strength"] >= 2.0
|
||||
|
||||
def test_strength_within_bounds(self):
|
||||
"""Strength should never exceed 2.5."""
|
||||
candles = [
|
||||
{"open": 98, "high": 105, "low": 97, "close": 102},
|
||||
{"open": 104, "high": 105, "low": 96, "close": 99},
|
||||
]
|
||||
result = detect_price_action_signal(candles)
|
||||
|
||||
assert result["strength"] <= 2.5
|
||||
assert result["strength"] >= 0.0
|
||||
|
||||
def test_proximity_boost_when_at_liquidity_level(self):
|
||||
"""Proximity to liquidity level should boost strength."""
|
||||
candles = [
|
||||
{"open": 98, "high": 105, "low": 97, "close": 102},
|
||||
{"open": 104, "high": 105, "low": 96, "close": 99},
|
||||
]
|
||||
liquidity_data = {"nearest_low": 99.5, "nearest_high": 105.5}
|
||||
|
||||
result = detect_price_action_signal(candles, liquidity_data=liquidity_data)
|
||||
|
||||
# When bearish engulfing near low liquidity level, proximity should be marked
|
||||
if result["pattern_type"] == "ENGULFING":
|
||||
# Just verify the field exists and is valid
|
||||
assert result["proximity_to_level"] in [None, "AT_LEVEL", "NEAR_LEVEL"]
|
||||
|
||||
|
||||
class TestLiquidityPriceActionIntegration:
|
||||
"""Integration tests: liquidity detection + price action + signal scoring."""
|
||||
|
||||
def test_liquidity_and_pa_together(self):
|
||||
"""When liquidity + PA both detect signals, should return valid data."""
|
||||
# Create realistic candles
|
||||
candles = [
|
||||
{"high": 100, "low": 90, "close": 95, "open": 92, "volume": 1000},
|
||||
{"high": 110, "low": 100, "close": 105, "open": 100, "volume": 1200}, # Swing high
|
||||
{"high": 105, "low": 95, "close": 100, "open": 105, "volume": 1100},
|
||||
{"high": 110, "low": 100, "close": 105, "open": 100, "volume": 1200},
|
||||
{"high": 104, "low": 95, "close": 96, "open": 104, "volume": 1300}, # Pin bar
|
||||
]
|
||||
|
||||
liquidity = detect_liquidity_levels(candles, pivot_lookback=3)
|
||||
pa = detect_price_action_signal(candles, liquidity_data=liquidity)
|
||||
|
||||
# Both should have valid outputs
|
||||
assert liquidity is not None
|
||||
assert pa is not None
|
||||
|
||||
def test_scoring_with_liquidity_and_pa(self):
|
||||
"""Score computation should handle liquidity + PA data."""
|
||||
indicators = {
|
||||
"close": [100, 102, 101, 103, 102],
|
||||
"sma_20": [100, 101, 101, 102, 102],
|
||||
"rsi_14": [50, 55, 52, 58, 55],
|
||||
"bollinger_bands": {
|
||||
"upper": [110, 111, 110, 112, 111],
|
||||
"middle": [100, 101, 101, 102, 102],
|
||||
"lower": [90, 91, 92, 92, 93],
|
||||
},
|
||||
"liquidity_levels": {
|
||||
"nearest_high": 110.0,
|
||||
"nearest_low": 90.0,
|
||||
},
|
||||
"pa_signal": {
|
||||
"pattern_type": "ENGULFING",
|
||||
"direction": "BULLISH",
|
||||
"strength": 2.0,
|
||||
"proximity_to_level": "AT_LEVEL",
|
||||
},
|
||||
}
|
||||
|
||||
_, _, _, score, raw_scores = _compute_adjusted_score(
|
||||
close_price=102,
|
||||
bb=indicators["bollinger_bands"],
|
||||
rsi=indicators.get("rsi_14"),
|
||||
sma=indicators.get("sma_20"),
|
||||
macd_data=None,
|
||||
st_data=None,
|
||||
vol_data=None,
|
||||
liquidity_data=indicators.get("liquidity_levels"),
|
||||
pa_signal=indicators.get("pa_signal"),
|
||||
)
|
||||
|
||||
# Score should be positive (BULLISH signals)
|
||||
assert score > 0 or score == 0 # Allow 0 for neutral
|
||||
# Raw scores should have liquidity_sweep and price_action_reversal
|
||||
assert "liquidity_sweep" in raw_scores
|
||||
assert "price_action_reversal" in raw_scores
|
||||
|
||||
|
||||
class TestEdgeCases:
|
||||
"""Edge case and error handling tests."""
|
||||
|
||||
def test_empty_candles_list(self):
|
||||
"""Should handle empty candles gracefully."""
|
||||
candles = []
|
||||
|
||||
liquidity = detect_liquidity_levels(candles)
|
||||
pa = detect_price_action_signal(candles)
|
||||
|
||||
assert liquidity is not None
|
||||
assert pa is not None
|
||||
assert pa["pattern_type"] is None
|
||||
|
||||
def test_zero_range_candles(self):
|
||||
"""Should handle candles with zero range (OHLC all equal)."""
|
||||
candles = [
|
||||
{"high": 100, "low": 100, "close": 100, "open": 100},
|
||||
{"high": 100, "low": 100, "close": 100, "open": 100},
|
||||
]
|
||||
|
||||
result = detect_price_action_signal(candles)
|
||||
|
||||
# Should not crash; pattern_type should be None (no pattern)
|
||||
assert result["pattern_type"] is None
|
||||
|
||||
def test_extremely_large_values(self):
|
||||
"""Should handle very large price values."""
|
||||
candles = [
|
||||
{"high": 1e6, "low": 9e5, "close": 9.5e5, "open": 9.2e5},
|
||||
{"high": 1.1e6, "low": 1e6, "close": 1.05e6, "open": 1e6},
|
||||
]
|
||||
|
||||
result = detect_price_action_signal(candles)
|
||||
|
||||
# Should handle large values without overflow
|
||||
assert result is not None
|
||||
assert result["strength"] <= 2.5
|
||||
|
||||
def test_extremely_small_values(self):
|
||||
"""Should handle very small price values."""
|
||||
candles = [
|
||||
{"high": 0.001, "low": 0.0009, "close": 0.00095, "open": 0.00092},
|
||||
{"high": 0.0011, "low": 0.001, "close": 0.00105, "open": 0.001},
|
||||
]
|
||||
|
||||
result = detect_price_action_signal(candles)
|
||||
|
||||
# Should handle small values without underflow/precision loss
|
||||
assert result is not None
|
||||
Reference in New Issue
Block a user