Files
Chan/strategies/MakerEdgeProbe.py
T
jackyu66gitandCursor 8ee11317d3 fix(web): 自动刷新保留 K 线视窗;威科夫与图表增量更新
自动刷新改用 tail update 与 scrollToPosition 恢复视窗,避免 setData 后跳到最右;拆分 chart_tv 模块并扩展 analyze/recent API。同步威科夫分析、pipeline 增量构建及相关策略与配置。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-25 22:57:43 +08:00

435 lines
16 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
MakerEdgeProbe — Freqtrade Dry-run 探针(过渡用)。
正式 Maker / L2 / Edge 采集已迁移到:
nautilus_mm/ NautilusTrader,独立 .venv
本策略仍可用于 Freqtrade 侧对照;新开发请走 nautilus_mm。
运行 Nautilus
cd nautilus_mm && ./scripts/run_probe.sh
分析:
cd nautilus_mm && ./scripts/analyze.sh
"""
from __future__ import annotations
import logging
import time
from datetime import datetime, timedelta, timezone
from typing import Optional
import numpy as np
import talib.abstract as ta
from pandas import DataFrame
from freqtrade.persistence import Trade, Order
from freqtrade.strategy import IStrategy
from maker_edge_logger import MakerEdgeLogger
logger = logging.getLogger(__name__)
class MakerEdgeProbe(IStrategy):
INTERFACE_VERSION = 3
timeframe = "1m"
can_short = True
process_only_new_candles = False
startup_candle_count = 60
minimal_roi = {"0": 0.01}
stoploss = -0.002
trailing_stop = False
use_exit_signal = False
order_types = {
"entry": "limit",
"exit": "limit",
"stoploss": "market",
"stoploss_on_exchange": False,
}
order_time_in_force = {"entry": "GTC", "exit": "GTC"}
tick_size = 0.1
quote_depth_ticks = 1
max_leverage = 1.0
stake_pct = 0.003
max_hold_minutes = 5
edge_exit_pct = 0.0002
adverse_exit_pct = 0.0008
cooldown_minutes = 5
ob_levels = 10
trade_lookback = 100
ema_slope_thr = 0.0002
book_sample_every_sec = 2.0
_logger: MakerEdgeLogger | None = None
_last_mid: float | None = None
_last_book_sample: float = 0.0
_last_entry_time: Optional[datetime] = None
_recent_high: float = 0.0
_recent_low: float = 0.0
_pending_quote_id: Optional[str] = None
_fill_by_trade: dict[int, str] = {}
def bot_start(self, **kwargs) -> None:
self._logger = MakerEdgeLogger(levels=self.ob_levels)
self._fill_by_trade = {}
logger.info("MakerEdgeProbe started. log_dir=%s", self._logger.log_dir)
def _get_logger(self) -> MakerEdgeLogger:
if self._logger is None:
self._logger = MakerEdgeLogger(levels=self.ob_levels)
return self._logger
def _fetch_trades(self, pair: str) -> list:
try:
ex = self.dp._exchange
if ex is None:
return []
api = getattr(ex, "_api", None) or getattr(ex, "api", None)
if api is None:
return []
return api.fetch_trades(pair, limit=self.trade_lookback) or []
except Exception as e:
logger.debug("fetch_trades failed: %s", e)
return []
def _inventory(self) -> float:
try:
inv = 0.0
for t in Trade.get_open_trades():
amt = float(t.amount or 0.0)
inv += -amt if t.is_short else amt
return inv
except Exception:
return 0.0
def _market_state(self, pair: str) -> dict:
state = {
"trend_state": "UNKNOWN",
"atr_pct": None,
"volatility_regime": "UNKNOWN",
"ema_slope": None,
}
try:
df, _ = self.dp.get_analyzed_dataframe(pair, self.timeframe)
if df is None or len(df) == 0:
return state
last = df.iloc[-1]
slope = float(last.get("ema_slope") or 0.0)
atr_pct = float(last.get("atr_pct") or 0.0)
state["ema_slope"] = slope
state["atr_pct"] = atr_pct
if bool(last.get("trend_block", False)):
state["trend_state"] = "TREND_UP" if slope > 0 else "TREND_DOWN"
else:
state["trend_state"] = "RANGE"
# 波动分位代理
if "atr_pct" in df.columns:
med = float(df["atr_pct"].tail(60).median() or 0)
if atr_pct > med * 1.8:
state["volatility_regime"] = "HIGH"
elif atr_pct < med * 0.7:
state["volatility_regime"] = "LOW"
else:
state["volatility_regime"] = "NORMAL"
except Exception:
pass
return state
def _snapshot(self, pair: str):
ob = self.dp.orderbook(pair, self.ob_levels)
trades = self._fetch_trades(pair)
snap = MakerEdgeLogger.snapshot_from_orderbook(
ob,
levels=self.ob_levels,
recent_trades=trades,
last_mid=self._last_mid,
liq_proxy_low=self._recent_low or None,
liq_proxy_high=self._recent_high or None,
)
if snap.mid:
self._last_mid = snap.mid
return snap
def bot_loop_start(self, current_time: datetime, **kwargs) -> None:
if self.dp.runmode.value not in ("live", "dry_run"):
return
pair = self.config["exchange"]["pair_whitelist"][0]
try:
snap = self._snapshot(pair)
tick = self.dp.ticker(pair) or {}
last = float(tick.get("last") or tick.get("close") or 0.0) or snap.mid
df, _ = self.dp.get_analyzed_dataframe(pair, self.timeframe)
if df is not None and len(df):
self._recent_high = float(df.iloc[-1].get("roll_high") or self._recent_high or last)
self._recent_low = float(df.iloc[-1].get("roll_low") or self._recent_low or last)
lg = self._get_logger()
now = time.time()
# 盘口历史(成交前5s恶化检测依赖此)
if now - self._last_book_sample >= self.book_sample_every_sec:
self._last_book_sample = now
lg.record_book(snap, now=now)
if last:
lg.update_paths(pair, last, now=now)
except Exception as e:
logger.warning("bot_loop_start probe error: %s", e)
def populate_indicators(self, dataframe: DataFrame, metadata: dict) -> DataFrame:
df = dataframe
close, high, low = df["close"], df["high"], df["low"]
volume = df["volume"].astype(float)
close_c = close.clip(lower=low, upper=high)
hl = (high - low).replace(0, np.nan)
buy_frac = ((close_c - low) / hl).fillna(0.5).clip(0, 1)
sell_vol = volume * (1.0 - buy_frac)
buy_vol = volume * buy_frac
df["sell_vol"] = sell_vol
df["buy_vol"] = buy_vol
df["delta"] = buy_vol - sell_vol
vol_ma = volume.rolling(20, min_periods=5).mean()
df["shock_sell"] = (sell_vol > vol_ma * 3) & (df["delta"] < 0)
df["shock_buy"] = (buy_vol > vol_ma * 3) & (df["delta"] > 0)
drop = (close.shift(3) - low).clip(lower=0) / close.shift(3)
up = (high - close.shift(3)).clip(lower=0) / close.shift(3)
df["de_sell"] = (sell_vol.rolling(3).sum() / (drop.replace(0, np.nan) * 1e4)).replace(
[np.inf, -np.inf], np.nan
).fillna(0)
df["de_buy"] = (buy_vol.rolling(3).sum() / (up.replace(0, np.nan) * 1e4)).replace(
[np.inf, -np.inf], np.nan
).fillna(0)
df["ema26"] = ta.EMA(df, timeperiod=26)
df["ema_slope"] = ((df["ema26"] - df["ema26"].shift(5)) / close).fillna(0)
df["trend_block"] = df["ema_slope"].abs() > self.ema_slope_thr
df["atr"] = ta.ATR(df, timeperiod=20)
df["atr_pct"] = (df["atr"] / close).fillna(0)
s_ma, s_ref = sell_vol.rolling(3).mean(), sell_vol.rolling(8).mean()
df["sell_exhaust"] = (s_ma < s_ref * 0.75) & (low >= low.rolling(8).min().shift(1))
b_ma, b_ref = buy_vol.rolling(3).mean(), buy_vol.rolling(8).mean()
df["buy_exhaust"] = (b_ma < b_ref * 0.75) & (high <= high.rolling(8).max().shift(1))
df["roll_high"] = high.rolling(60, min_periods=10).max()
df["roll_low"] = low.rolling(60, min_periods=10).min()
return df
def populate_entry_trend(self, dataframe: DataFrame, metadata: dict) -> DataFrame:
df = dataframe
long_c = (
(~df["trend_block"])
& df["shock_sell"].rolling(5).max().astype(bool)
& (df["de_sell"] > 10)
& df["sell_exhaust"]
)
short_c = (
(~df["trend_block"])
& df["shock_buy"].rolling(5).max().astype(bool)
& (df["de_buy"] > 10)
& df["buy_exhaust"]
)
df.loc[long_c, ["enter_long", "enter_tag"]] = (1, "probe_bid_lp")
df.loc[short_c, ["enter_short", "enter_tag"]] = (1, "probe_ask_lp")
return df
def populate_exit_trend(self, dataframe: DataFrame, metadata: dict) -> DataFrame:
dataframe["exit_long"] = 0
dataframe["exit_short"] = 0
return dataframe
def custom_entry_price(
self,
pair: str,
trade: Trade | None,
current_time: datetime,
proposed_rate: float,
entry_tag: str | None,
side: str,
**kwargs,
) -> float:
offset = self.quote_depth_ticks * self.tick_size
try:
snap = self._snapshot(pair)
price = snap.best_bid - offset if side == "long" else snap.best_ask + offset
state = self._market_state(pair)
qid = self._get_logger().create_quote(
pair=pair,
side="bid" if side == "long" else "ask",
quote_price=price,
inventory=self._inventory(),
snap=snap,
reason=entry_tag or "entry",
trade_id=trade.id if trade else None,
state=state,
)
self._pending_quote_id = qid
return price
except Exception as e:
logger.debug("custom_entry_price: %s", e)
return proposed_rate - offset if side == "long" else proposed_rate + offset
def confirm_trade_entry(
self,
pair: str,
order_type: str,
amount: float,
rate: float,
time_in_force: str,
current_time: datetime,
entry_tag: str | None,
side: str,
**kwargs,
) -> bool:
if self._last_entry_time:
last = self._last_entry_time
if last.tzinfo is None:
last = last.replace(tzinfo=timezone.utc)
now = current_time if current_time.tzinfo else current_time.replace(tzinfo=timezone.utc)
if now - last < timedelta(minutes=self.cooldown_minutes):
return False
try:
df, _ = self.dp.get_analyzed_dataframe(pair, self.timeframe)
if df is not None and len(df) and bool(df.iloc[-1].get("trend_block", False)):
return False
snap = self._snapshot(pair)
if side == "long" and snap.bid_depth_1 < snap.ask_depth_1 * 0.7:
return False
if side == "short" and snap.ask_depth_1 < snap.bid_depth_1 * 0.7:
return False
except Exception:
pass
self._last_entry_time = current_time
return True
def check_entry_timeout(
self, pair: str, trade: Trade, order: Order, current_time: datetime, **kwargs
) -> bool:
"""超时撤单 → 记录 quote_cancel(坏时间未成交 vs 被动成交的对照)。"""
try:
snap = self._snapshot(pair)
self._get_logger().cancel_quote(
quote_id=self._pending_quote_id,
trade_id=trade.id,
reason="entry_timeout",
snap=snap,
)
except Exception as e:
logger.debug("cancel_quote on timeout: %s", e)
# False = 不额外强制取消;交给 unfilledtimeout 配置。若要立刻取消返回 True
return False
def order_filled(
self,
pair: str,
trade: Trade,
order: Order,
current_time: datetime,
**kwargs,
) -> None:
try:
lg = self._get_logger()
# 入场成交
if order.ft_order_side == trade.entry_side:
snap = self._snapshot(pair)
side = "short" if trade.is_short else "long"
# 粗分 fill_reasontime_to_fill 在 logger 内算;这里标 maker_hit
# 若成交前5s盘口已恶化 → toxic_passive 候选
det = lg.book_deterioration(side)
fill_reason = "toxic_passive" if det.get("pre_5s_deteriorated") else "maker_hit"
if self._pending_quote_id:
lg.bind_trade(self._pending_quote_id, trade.id)
fill_id = lg.log_fill(
pair=pair,
side=side,
fill_price=float(order.safe_price or trade.open_rate),
amount=float(order.safe_filled or order.safe_amount or 0),
inventory=self._inventory(),
snap=snap,
order_type=str(getattr(order, "order_type", None) or "limit"),
quote_id=self._pending_quote_id,
trade_id=trade.id,
fill_reason=fill_reason,
state=self._market_state(pair),
extra={"entry_tag": trade.enter_tag},
)
self._fill_by_trade[trade.id] = fill_id
self._pending_quote_id = None
else:
# 出场:把 exit_reason 挂到入场 fill,供 H2
fill_id = self._fill_by_trade.get(trade.id)
reason = trade.exit_reason or getattr(order, "ft_order_tag", None) or "exit"
if fill_id:
lg.attach_exit_reason(fill_id, str(reason))
except Exception as e:
logger.warning("order_filled log error: %s", e)
def custom_exit(
self,
pair: str,
trade: Trade,
current_time: datetime,
current_rate: float,
current_profit: float,
**kwargs,
):
open_time = trade.open_date_utc
if open_time.tzinfo is None:
open_time = open_time.replace(tzinfo=timezone.utc)
now = current_time if current_time.tzinfo else current_time.replace(tzinfo=timezone.utc)
if now - open_time >= timedelta(minutes=self.max_hold_minutes):
return "probe_time"
entry = trade.open_rate
edge = (
(current_rate - entry) / entry
if not trade.is_short
else (entry - current_rate) / entry
)
if edge >= self.edge_exit_pct:
return "probe_edge_restore"
if edge <= -self.adverse_exit_pct:
return "probe_adverse"
# 趋势切换 → 撤流动性思维
try:
st = self._market_state(pair)
if st.get("trend_state") in ("TREND_UP", "TREND_DOWN"):
# 持仓方向与趋势相反时更危险
if (not trade.is_short and st["trend_state"] == "TREND_DOWN") or (
trade.is_short and st["trend_state"] == "TREND_UP"
):
return "probe_trend_cancel"
except Exception:
pass
return None
def leverage(
self, pair: str, current_time: datetime, current_rate: float,
proposed_leverage: float, max_leverage: float, entry_tag: Optional[str],
side: str, **kwargs,
) -> float:
return min(self.max_leverage, float(max_leverage))
def custom_stake_amount(
self, pair: str, current_time: datetime, current_rate: float,
proposed_stake: float, min_stake: Optional[float], max_stake: float,
leverage: float, entry_tag: Optional[str], side: str, **kwargs,
) -> float:
try:
if self.wallets:
free = self.wallets.get_free(self.config["stake_currency"])
stake = free * self.stake_pct
if min_stake:
stake = max(stake, min_stake)
return min(stake, max_stake)
except Exception:
pass
return min(proposed_stake * self.stake_pct, max_stake)