diff --git a/backend/app/api/settings.py b/backend/app/api/settings.py index 3a2c49b..9ea9ce7 100644 --- a/backend/app/api/settings.py +++ b/backend/app/api/settings.py @@ -14,10 +14,12 @@ router = APIRouter(prefix="/api/settings", tags=["settings"]) KEYS = ( "fee_rate", - "exit_move_points", + "exit_move_pct", "rest_seconds", - "max_rounds", "initial_equity", + "leverage", + "min_option_hours", + "min_option_leverage", "perp_qty_eth", "option_qty_eth", ) @@ -25,10 +27,12 @@ KEYS = ( class StrategySettingsBody(BaseModel): fee_rate: float | None = Field(default=None, ge=0, le=0.05) - exit_move_points: float | None = Field(default=None, ge=1, le=500) + exit_move_pct: float | None = Field(default=None, ge=0.1, le=50) rest_seconds: int | None = Field(default=None, ge=0, le=3600) - max_rounds: int | None = Field(default=None, ge=1, le=20) initial_equity: float | None = Field(default=None, ge=1000) + leverage: float | None = Field(default=None, ge=1, le=125) + min_option_hours: float | None = Field(default=None, ge=1, le=720) + min_option_leverage: float | None = Field(default=None, ge=1, le=10000) perp_qty_eth: float | None = Field(default=None, ge=0.01, le=100) option_qty_eth: float | None = Field(default=None, ge=0.01, le=100) @@ -38,18 +42,24 @@ def _read_settings() -> dict: s = get_settings() return { "fee_rate": float(db.get_setting("fee_rate", str(s.fee_rate)) or s.fee_rate), - "exit_move_points": float( - db.get_setting("exit_move_points", str(s.exit_move_points)) or s.exit_move_points + "exit_move_pct": float( + db.get_setting("exit_move_pct", str(s.exit_move_pct)) or s.exit_move_pct ), "rest_seconds": int( float(db.get_setting("rest_seconds", str(s.rest_seconds)) or s.rest_seconds) ), - "max_rounds": int( - float(db.get_setting("max_rounds", str(s.max_rounds)) or s.max_rounds) - ), "initial_equity": float( db.get_setting("initial_equity", str(s.initial_equity)) or s.initial_equity ), + "leverage": float(db.get_setting("leverage", str(s.leverage)) or s.leverage), + "min_option_hours": float( + db.get_setting("min_option_hours", str(s.min_option_hours)) + or s.min_option_hours + ), + "min_option_leverage": float( + db.get_setting("min_option_leverage", str(s.min_option_leverage)) + or s.min_option_leverage + ), "perp_qty_eth": float( db.get_setting("perp_qty_eth", str(s.perp_qty_eth)) or s.perp_qty_eth ), diff --git a/backend/app/api/sim.py b/backend/app/api/sim.py index 9655ef4..42536b2 100644 --- a/backend/app/api/sim.py +++ b/backend/app/api/sim.py @@ -7,10 +7,10 @@ from pydantic import BaseModel, Field from ..market import get_gateway from ..models.db import get_db +from ..sim.ledger import Ledger from ..sim.matcher import Matcher from ..strategy.clock import window_key from ..strategy.group import next_group_id -from ..strategy.signal import decide from .auth import require_user router = APIRouter(prefix="/api/sim", tags=["sim"]) @@ -38,37 +38,45 @@ async def sim_open_group( body: ManualOpenBody | None = None, ) -> dict: gw = get_gateway() - # 开仓前按现价强制重选 ATM,避免沿用启动时的旧行权价 - try: - await gw.ensure_atm_async(force=True) - except Exception as e: - raise HTTPException(status_code=503, detail=f"ATM 对齐失败: {e}") from e - snap = gw.snapshot() - if not snap.pair or not snap.call or not snap.put: - raise HTTPException(status_code=503, detail="行情未就绪") + pick = await gw.pick_for_open_async() + if pick is None: + raise HTTPException( + status_code=409, + detail="无合格期权:请检查剩余时长(≥设置小时)与杠杆(现价/卖一)", + ) + force = (body.force_option_side if body else None) or None if force in ("call", "put"): option_side = force perp_side = "short" if force == "call" else "long" bias = "manual_" + force + option_ask = pick.call_ask if force == "call" else pick.put_ask + from ..config import get_settings + from ..market.instruments import option_leverage + from ..sim.ledger import Ledger as Led + + s = get_settings() + min_lev = Led().get_setting_float("min_option_leverage", s.min_option_leverage) + lev = option_leverage(pick.underlying_px, option_ask) + if lev is None or lev < min_lev: + raise HTTPException( + status_code=409, + detail=f"强制方向杠杆不足: {lev or 0:.1f} < {min_lev:.0f}", + ) else: - sig = decide(snap.call.ask, snap.put.ask) - if sig is None: - raise HTTPException(status_code=409, detail="Call/Put 卖一相等,跳过") - option_side = sig.option_side - perp_side = sig.perp_side - bias = sig.bias + option_side = pick.option_side + perp_side = pick.perp_side + bias = pick.bias option_inst = ( - snap.pair.call_inst_id if option_side == "call" else snap.pair.put_inst_id + pick.pair.call_inst_id if option_side == "call" else pick.pair.put_inst_id ) - entry_idx = snap.index_px or (snap.perp.mark_px if snap.perp else None) - if entry_idx is None: - raise HTTPException(status_code=503, detail="无指数/标记价") wkey = window_key() db = get_db() - count = len(db.fetchall("SELECT group_id FROM groups WHERE group_id LIKE ?", (f"G-{wkey}-%",))) + count = len( + db.fetchall("SELECT group_id FROM groups WHERE group_id LIKE ?", (f"G-{wkey}-%",)) + ) gid = next_group_id(count) r = Matcher().open_group( group_id=gid, @@ -76,13 +84,21 @@ async def sim_open_group( option_side=option_side, perp_side=perp_side, option_inst_id=option_inst, - entry_index_px=float(entry_idx), - strike=snap.pair.strike, - expiry_ymd=snap.pair.expiry_ymd, + entry_index_px=float(pick.underlying_px), + strike=pick.pair.strike, + expiry_ymd=pick.pair.expiry_ymd, ) if not r.ok: raise HTTPException(status_code=400, detail=r.detail) - return {"ok": True, **(r.data or {}), "detail": r.detail} + return { + "ok": True, + **(r.data or {}), + "detail": r.detail, + "option_leverage": pick.option_leverage, + "hours_left": pick.hours_left, + "expiry_ymd": pick.pair.expiry_ymd, + "strike": pick.pair.strike, + } @router.post("/close-group") diff --git a/backend/app/config.py b/backend/app/config.py index 32a885e..625e743 100644 --- a/backend/app/config.py +++ b/backend/app/config.py @@ -38,11 +38,15 @@ class Settings(BaseSettings): fee_rate: float = 0.0005 initial_equity: float = 100_000.0 - max_rounds: int = 3 - open_hhmm: str = "16:00" - stop_open_hhmm: str = "08:00" - exit_move_points: float = 30.0 + max_rounds: int = 3 # 已不再强管控,仅兼容旧字段 + open_hhmm: str = "16:00" # 已废弃开仓窗 + stop_open_hhmm: str = "08:00" # 已废弃开仓窗 + exit_move_points: float = 30.0 # 旧字段,改用 exit_move_pct + exit_move_pct: float = 2.0 # 相对开仓指数波动 % 全平 rest_seconds: int = 300 + leverage: float = 3.0 # 永续杠杆 + min_option_hours: float = 12.0 # 期权最小剩余小时 + min_option_leverage: float = 100.0 # 现价/卖一权利金 下限 perp_qty_eth: float = 1.0 option_qty_eth: float = 2.0 option_ct_mult_default: float = 0.01 diff --git a/backend/app/market/gateway.py b/backend/app/market/gateway.py index f4854aa..7e29b68 100644 --- a/backend/app/market/gateway.py +++ b/backend/app/market/gateway.py @@ -4,18 +4,24 @@ from __future__ import annotations import asyncio import logging +from dataclasses import dataclass from typing import Any from ..config import Settings, get_settings from .book_cache import BookCache -from .instruments import next_session_expiry_ymd, select_option_pair +from .instruments import ( + hours_until_expiry, + list_eligible_expiry_ymds, + option_leverage, + select_option_pair, +) from .okx_rest import OkxRestClient from .okx_ws import OkxPublicWs from .types import MarketSnapshot, OptionPair logger = logging.getLogger(__name__) -# 现价偏离当前行权价超过该点数则重选 ATM(ETH 期权常见步进 5) +# 展示用:现价偏离当前行权超过该点数则重选 ATM(空仓) _ATM_DRIFT_POINTS = 5.0 @@ -29,6 +35,40 @@ def _has_open_position() -> bool: return False +def _strategy_floats() -> tuple[float, float]: + """(min_option_hours, min_option_leverage)""" + s = get_settings() + try: + from ..models.db import get_db + + db = get_db() + hours = float( + db.get_setting("min_option_hours", str(s.min_option_hours)) + or s.min_option_hours + ) + lev = float( + db.get_setting("min_option_leverage", str(s.min_option_leverage)) + or s.min_option_leverage + ) + return hours, lev + except Exception: + return s.min_option_hours, s.min_option_leverage + + +@dataclass(slots=True) +class OpenPick: + pair: OptionPair + option_side: str + perp_side: str + bias: str + call_ask: float + put_ask: float + option_ask: float + option_leverage: float + hours_left: float + underlying_px: float + + class MarketGateway: def __init__(self, settings: Settings | None = None) -> None: self.settings = settings or get_settings() @@ -64,25 +104,13 @@ class MarketGateway: await self.ws.stop() self.rest.close() - def align_instruments(self) -> OptionPair | None: - """同步:拉期权列表,选次日到期 ATM Call/Put,REST 预热盘口,切换 WS 订阅。""" + def _apply_pair(self, pair: OptionPair, *, mark: float, idx: float | None) -> OptionPair: s = self.settings - idx = self.rest.fetch_index_ticker(s.index_inst_id) - mark = self.rest.fetch_mark_price(s.perp_inst_id) or idx - if mark is None or mark <= 0: - raise RuntimeError("无法获取 ETH 标记/指数价格,无法选 ATM") - - instruments = self.rest.fetch_option_instruments(s.option_inst_family) - ymd = next_session_expiry_ymd() - pair = select_option_pair(instruments, mark_px=float(mark), expiry_ymd=ymd) - if pair is None: - raise RuntimeError(f"未找到到期 {ymd} 的 ATM Call/Put 合约 pair (family={s.option_inst_family})") - self._pair = pair self.cache.set_pair(pair) - self.cache.set_index_px(idx) + if idx is not None: + self.cache.set_index_px(idx) - # REST 预热:永续 + Call + Put for inst in (s.perp_inst_id, pair.call_inst_id, pair.put_inst_id): bids, asks, ts = self.rest.fetch_books(inst, sz=5) self.cache.upsert_book(inst, bids=bids, asks=asks, ts_ms=ts) @@ -94,15 +122,102 @@ class MarketGateway: self.cache.drop_except(keep) self.ws.set_instruments([s.perp_inst_id, pair.call_inst_id, pair.put_inst_id]) logger.info( - "aligned pair expiry=%s strike=%s call=%s put=%s mark=%.2f", + "aligned pair expiry=%s strike=%s call=%s put=%s mark=%.2f hours=%.1f", pair.expiry_ymd, pair.strike, pair.call_inst_id, pair.put_inst_id, mark, + hours_until_expiry(pair.expiry_ymd), ) return pair + def align_instruments(self) -> OptionPair | None: + """空仓展示:选剩余时长合格的最近到期 ATM(不校验期权杠杆)。""" + s = self.settings + idx = self.rest.fetch_index_ticker(s.index_inst_id) + mark = self.rest.fetch_mark_price(s.perp_inst_id) or idx + if mark is None or mark <= 0: + raise RuntimeError("无法获取 ETH 标记/指数价格,无法选 ATM") + + min_hours, _ = _strategy_floats() + instruments = self.rest.fetch_option_instruments(s.option_inst_family) + pair = select_option_pair( + instruments, mark_px=float(mark), min_hours=min_hours + ) + if pair is None: + raise RuntimeError( + f"未找到剩余≥{min_hours}h 的 ATM Call/Put (family={s.option_inst_family})" + ) + return self._apply_pair(pair, mark=float(mark), idx=idx) + + def pick_for_open(self) -> OpenPick | None: + """ + 开仓选约: + 1) 剩余时长 ≥ min_hours 的到期日(由近到远) + 2) 该到期 ATM 平值 + 3) 卖一比价定方向后校验 现价/卖一 ≥ min_option_leverage + """ + s = self.settings + min_hours, min_lev = _strategy_floats() + idx = self.rest.fetch_index_ticker(s.index_inst_id) + mark = self.rest.fetch_mark_price(s.perp_inst_id) or idx + if mark is None or mark <= 0: + return None + underlying = float(mark) + instruments = self.rest.fetch_option_instruments(s.option_inst_family) + eligible = list_eligible_expiry_ymds(instruments, min_hours=min_hours) + if not eligible: + logger.info("no expiry with hours>=%.1f", min_hours) + return None + + from ..strategy.signal import decide + + for ymd in eligible: + pair = select_option_pair( + instruments, mark_px=underlying, expiry_ymd=ymd + ) + if pair is None: + continue + call_bids, call_asks, _ = self.rest.fetch_books(pair.call_inst_id, sz=5) + put_bids, put_asks, _ = self.rest.fetch_books(pair.put_inst_id, sz=5) + call_ask = call_asks[0].px if call_asks else None + put_ask = put_asks[0].px if put_asks else None + sig = decide(call_ask, put_ask) + if sig is None: + continue + opt_ask = sig.call_ask if sig.option_side == "call" else sig.put_ask + lev = option_leverage(underlying, opt_ask) + hours_left = hours_until_expiry(ymd) + if lev is None or lev + 1e-9 < min_lev: + logger.info( + "skip expiry=%s strike=%.0f side=%s lev=%s need>=%.0f hours=%.1f", + ymd, + pair.strike, + sig.option_side, + f"{lev:.1f}" if lev else "n/a", + min_lev, + hours_left, + ) + continue + self._apply_pair(pair, mark=underlying, idx=idx) + # 写入刚拉的盘口,避免 WS 尚未推送 + self.cache.upsert_book(pair.call_inst_id, bids=call_bids, asks=call_asks) + self.cache.upsert_book(pair.put_inst_id, bids=put_bids, asks=put_asks) + return OpenPick( + pair=pair, + option_side=sig.option_side, + perp_side=sig.perp_side, + bias=sig.bias, + call_ask=float(sig.call_ask), + put_ask=float(sig.put_ask), + option_ask=float(opt_ask), + option_leverage=float(lev), + hours_left=hours_left, + underlying_px=underlying, + ) + return None + async def realign_async(self) -> OptionPair | None: old = self._pair pair = await asyncio.to_thread(self.align_instruments) @@ -122,6 +237,23 @@ class MarketGateway: ) return pair + async def pick_for_open_async(self) -> OpenPick | None: + old = self._pair + pick = await asyncio.to_thread(self.pick_for_open) + if pick and ( + old is None + or pick.pair.call_inst_id != old.call_inst_id + or pick.pair.put_inst_id != old.put_inst_id + ): + await self.ws.resubscribe( + [ + self.settings.perp_inst_id, + pick.pair.call_inst_id, + pick.pair.put_inst_id, + ] + ) + return pick + def _mark_for_atm(self) -> float | None: snap = self.snapshot() if snap.perp and snap.perp.mark_px: @@ -133,11 +265,10 @@ class MarketGateway: return None def atm_needs_realign(self, mark_px: float | None = None) -> bool: - """到期日变了,或现价已偏离当前行权价超过阈值。""" if self._pair is None: return True - want = next_session_expiry_ymd() - if self._pair.expiry_ymd != want: + min_hours, _ = _strategy_floats() + if hours_until_expiry(self._pair.expiry_ymd) + 1e-9 < min_hours: return True mark = mark_px if mark_px is not None else self._mark_for_atm() if mark is None or mark <= 0: @@ -145,14 +276,15 @@ class MarketGateway: return abs(float(self._pair.strike) - float(mark)) >= _ATM_DRIFT_POINTS async def ensure_atm_async(self, *, force: bool = False) -> OptionPair | None: - """空仓时按现价对齐 ATM。有持仓时不切换,避免盯市合约被换掉。""" + """空仓时按剩余时长+ATM 对齐。有持仓不切换。""" if _has_open_position(): return self._pair if force or self.atm_needs_realign(): logger.info( - "ATM realign force=%s old_strike=%s", + "ATM realign force=%s old_strike=%s old_exp=%s", force, self._pair.strike if self._pair else None, + self._pair.expiry_ymd if self._pair else None, ) return await self.realign_async() return self._pair @@ -164,7 +296,6 @@ class MarketGateway: return self.snapshot().to_dict() async def _refresh_loop(self) -> None: - """周期性刷新指数价;空仓时按到期/ATM 偏离重对齐。""" while True: await asyncio.sleep(30) try: @@ -184,7 +315,6 @@ class MarketGateway: logger.warning("market refresh failed: %s", e) -# 进程级单例(FastAPI lifespan 注入) _gateway: MarketGateway | None = None diff --git a/backend/app/market/instruments.py b/backend/app/market/instruments.py index ddc5d06..91719ae 100644 --- a/backend/app/market/instruments.py +++ b/backend/app/market/instruments.py @@ -1,4 +1,4 @@ -"""合约选择:次日 16:00(上海)到期 + ATM 行权价(暂定默认,待拍板可改)。""" +"""合约选择:剩余时长过滤 + ATM 平值期权。""" from __future__ import annotations @@ -42,13 +42,15 @@ def expiry_ms_from_ymd(ymd: str) -> int: return int(dt.timestamp() * 1000) +def hours_until_expiry(ymd: str, now: datetime | None = None) -> float: + """距到期剩余小时(可为负)。""" + n = (now or datetime.now(tz=_SH)).astimezone(_SH) + left_ms = expiry_ms_from_ymd(ymd) - int(n.timestamp() * 1000) + return left_ms / 3_600_000.0 + + def next_session_expiry_ymd(now: datetime | None = None) -> str: - """ - 业务约定:开仓选「次日 16:00」到期。 - - 上海时间 >= 当日 16:00:目标到期日 = 次日 - - 上海时间 < 当日 16:00:目标到期日 = 当日(当日 16:00 到期仍可用作盘口对齐/预热) - 正式开仓窗从当日 16:00 起,届时「次日」即日历次日。 - """ + """兼容旧逻辑:次日/当日 16:00 到期键(展示/测试用)。""" now_sh = (now or datetime.now(tz=_SH)).astimezone(_SH) open_today = now_sh.replace(hour=16, minute=0, second=0, microsecond=0) if now_sh >= open_today: @@ -64,30 +66,19 @@ def pick_atm_strike(strikes: list[float], mark_px: float) -> float | None: return min(strikes, key=lambda s: (abs(s - mark_px), s)) -def select_option_pair( +def _complete_by_expiry( instruments: list[dict[str, Any]], - *, - mark_px: float, - expiry_ymd: str | None = None, - now: datetime | None = None, -) -> OptionPair | None: - """ - 从 live 合约列表中选出:目标到期日 + ATM 同行权价 Call/Put。 - 行权价规则暂定 ATM(最接近标记/指数价);待拍板后可替换。 - """ - ymd = expiry_ymd or next_session_expiry_ymd(now) - by_strike: dict[float, dict[str, str]] = {} - +) -> dict[str, dict[float, dict[str, str]]]: + """expiry_ymd -> strike -> {C|P: instId},仅完整 Call+Put。""" + by_exp: dict[str, dict[float, dict[str, str]]] = {} for row in instruments: if not isinstance(row, dict): continue state = str(row.get("state") or "live").lower() if state and state != "live": continue - inst_id = str(row.get("instId") or "") y, stk, opt = parse_option_inst_id(inst_id) - if y is None or stk is None or opt is None: exp = safe_float(row.get("expTime")) if exp: @@ -96,20 +87,77 @@ def select_option_pair( stk = safe_float(row.get("stk")) opt_raw = str(row.get("optType") or "").upper() opt = opt_raw if opt_raw in ("C", "P") else None - - if not inst_id or y != ymd or stk is None or opt not in ("C", "P"): + if not inst_id or not y or stk is None or opt not in ("C", "P"): continue - by_strike.setdefault(float(stk), {})[opt] = inst_id + by_exp.setdefault(y, {}).setdefault(float(stk), {})[opt] = inst_id - complete = {s: v for s, v in by_strike.items() if "C" in v and "P" in v} + out: dict[str, dict[float, dict[str, str]]] = {} + for ymd, strikes in by_exp.items(): + complete = {s: v for s, v in strikes.items() if "C" in v and "P" in v} + if complete: + out[ymd] = complete + return out + + +def list_eligible_expiry_ymds( + instruments: list[dict[str, Any]], + *, + min_hours: float, + now: datetime | None = None, +) -> list[str]: + """剩余时间 >= min_hours 的到期日,由近到远。""" + complete = _complete_by_expiry(instruments) + eligible = [ + ymd + for ymd in complete + if hours_until_expiry(ymd, now) + 1e-9 >= float(min_hours) + ] + return sorted(eligible, key=lambda y: expiry_ms_from_ymd(y)) + + +def select_option_pair( + instruments: list[dict[str, Any]], + *, + mark_px: float, + expiry_ymd: str | None = None, + min_hours: float | None = None, + now: datetime | None = None, +) -> OptionPair | None: + """ + 选 ATM Call/Put。 + - 若给 expiry_ymd:在该到期日选平值。 + - 若给 min_hours:选「剩余时长合格」中最近到期日的平值。 + - 否则回退 next_session_expiry_ymd。 + """ + complete = _complete_by_expiry(instruments) if not complete: return None - atm = pick_atm_strike(list(complete.keys()), mark_px) + if expiry_ymd: + ymd = expiry_ymd + if ymd not in complete: + return None + elif min_hours is not None: + eligible = list_eligible_expiry_ymds( + instruments, min_hours=min_hours, now=now + ) + if not eligible: + return None + ymd = eligible[0] + else: + ymd = next_session_expiry_ymd(now) + if ymd not in complete: + # 回退到最近合格到期 + eligible = list_eligible_expiry_ymds(instruments, min_hours=0, now=now) + if not eligible: + return None + ymd = eligible[0] + + strikes_map = complete[ymd] + atm = pick_atm_strike(list(strikes_map.keys()), mark_px) if atm is None: return None - - legs = complete[atm] + legs = strikes_map[atm] return OptionPair( expiry_ymd=ymd, expiry_ms=expiry_ms_from_ymd(ymd), @@ -117,3 +165,10 @@ def select_option_pair( call_inst_id=legs["C"], put_inst_id=legs["P"], ) + + +def option_leverage(underlying_px: float, premium_ask: float) -> float | None: + """现价 / 卖一权利金。""" + if underlying_px <= 0 or premium_ask is None or premium_ask <= 0: + return None + return float(underlying_px) / float(premium_ask) diff --git a/backend/app/models/db.py b/backend/app/models/db.py index cbbd91e..23fdf67 100644 --- a/backend/app/models/db.py +++ b/backend/app/models/db.py @@ -149,8 +149,12 @@ class Database: "fee_rate": str(s.fee_rate), "initial_equity": str(s.initial_equity), "exit_move_points": str(s.exit_move_points), + "exit_move_pct": str(s.exit_move_pct), "rest_seconds": str(s.rest_seconds), "max_rounds": str(s.max_rounds), + "leverage": str(s.leverage), + "min_option_hours": str(s.min_option_hours), + "min_option_leverage": str(s.min_option_leverage), "perp_qty_eth": str(s.perp_qty_eth), "option_qty_eth": str(s.option_qty_eth), } diff --git a/backend/app/sim/matcher.py b/backend/app/sim/matcher.py index 59c0dfa..c5fea9a 100644 --- a/backend/app/sim/matcher.py +++ b/backend/app/sim/matcher.py @@ -385,6 +385,7 @@ class Matcher: "option_upl": 0.0, "index_px": None, "move_points": 0.0, + "move_pct": 0.0, "premium_gap": None, } gw = get_gateway() @@ -428,8 +429,12 @@ class Matcher: entry_idx = float(pos["entry_index_px"] or 0) move = abs(float(index_px) - entry_idx) if index_px is not None and entry_idx else 0.0 + move_pct = (move / entry_idx * 100.0) if entry_idx > 0 else 0.0 initial_premium = float(pos["initial_premium"] or 0) premium_gap = initial_premium - perp_upl + leverage = self.ledger.get_setting_float("leverage", s.leverage) + notional = abs(perp_entry * perp_qty) + margin = notional / leverage if leverage > 0 else None group_id = pos.get("group_id") g = ( @@ -454,6 +459,9 @@ class Matcher: "perp_entry_px": perp_entry, "perp_qty_eth": perp_qty, "perp_mark_px": float(mark) if mark is not None else None, + "perp_notional": notional, + "perp_margin": margin, + "leverage": leverage, "option_inst_id": pos.get("option_inst_id"), "option_entry_px": float(pos["option_entry_px"] or 0), "option_qty_eth": float(pos["option_qty_eth"] or 0), @@ -466,6 +474,7 @@ class Matcher: "index_px": index_px, "entry_index_px": entry_idx, "move_points": move, + "move_pct": move_pct, "initial_premium": initial_premium, "premium_gap": premium_gap, "status": pos.get("status"), diff --git a/backend/app/strategy/__init__.py b/backend/app/strategy/__init__.py index cd50b46..3821cf9 100644 --- a/backend/app/strategy/__init__.py +++ b/backend/app/strategy/__init__.py @@ -1,5 +1,4 @@ from .clock import can_open_new, window_key -from .engine import StrategyEngine, get_engine, set_engine from .exits import check_exits from .group import next_group_id from .signal import Signal, decide @@ -15,3 +14,15 @@ __all__ = [ "set_engine", "window_key", ] + + +def __getattr__(name: str): + if name in ("StrategyEngine", "get_engine", "set_engine"): + from .engine import StrategyEngine, get_engine, set_engine + + return { + "StrategyEngine": StrategyEngine, + "get_engine": get_engine, + "set_engine": set_engine, + }[name] + raise AttributeError(f"module {__name__!r} has no attribute {name!r}") diff --git a/backend/app/strategy/clock.py b/backend/app/strategy/clock.py index 3efc6e0..3e8fa3a 100644 --- a/backend/app/strategy/clock.py +++ b/backend/app/strategy/clock.py @@ -1,8 +1,8 @@ -"""业务窗时钟:16:00 开 → 08:00 停开;轮次与休息。""" +"""日历日分组键(开仓时间窗已取消,由期权剩余时长约束)。""" from __future__ import annotations -from datetime import datetime, timedelta +from datetime import datetime from zoneinfo import ZoneInfo _SH = ZoneInfo("Asia/Shanghai") @@ -12,28 +12,9 @@ def now_sh(now: datetime | None = None) -> datetime: return (now or datetime.now(tz=_SH)).astimezone(_SH) -def parse_hhmm(s: str) -> tuple[int, int]: - parts = (s or "16:00").strip().split(":") - return int(parts[0]), int(parts[1]) if len(parts) > 1 else 0 - - def window_key(now: datetime | None = None) -> str: - """ - 业务窗键:若当前 >= 当日 16:00,窗从今日 16:00 起,键=今日日期; - 若 < 16:00,仍可能属于「昨日起的窗」(到今日 08:00),键=昨日。 - """ - n = now_sh(now) - open_h, open_m = 16, 0 - stop_h, stop_m = 8, 0 - today_open = n.replace(hour=open_h, minute=open_m, second=0, microsecond=0) - today_stop = n.replace(hour=stop_h, minute=stop_m, second=0, microsecond=0) - if n >= today_open: - return n.strftime("%Y%m%d") - if n < today_stop: - # 仍在昨 16:00 开启的窗内 - return (n.date() - timedelta(days=1)).strftime("%Y%m%d") - # 08:00~16:00:不在开仓窗,键用「即将开始」的今日窗 - return n.strftime("%Y%m%d") + """组号日期键:日历日 YYYYMMDD。""" + return now_sh(now).strftime("%Y%m%d") def can_open_new( @@ -42,16 +23,8 @@ def can_open_new( open_hhmm: str = "16:00", stop_hhmm: str = "08:00", ) -> bool: - n = now_sh(now) - oh, om = parse_hhmm(open_hhmm) - sh, sm = parse_hhmm(stop_hhmm) - today_open = n.replace(hour=oh, minute=om, second=0, microsecond=0) - today_stop = n.replace(hour=sh, minute=sm, second=0, microsecond=0) - if n >= today_open: - return True - if n < today_stop: - return True - return False + """开仓窗已取消,始终允许(仍受期权剩余时长/杠杆筛选)。""" + return True def group_date_ymd(now: datetime | None = None) -> str: diff --git a/backend/app/strategy/engine.py b/backend/app/strategy/engine.py index 3743267..f73945c 100644 --- a/backend/app/strategy/engine.py +++ b/backend/app/strategy/engine.py @@ -1,4 +1,4 @@ -"""策略状态机:选向开仓 / 盯盘平仓 / 休息 / 限轮。""" +"""策略状态机:选向开仓 / 盯盘平仓 / 休息(无开仓窗、无轮次上限)。""" from __future__ import annotations @@ -12,10 +12,9 @@ from ..market import get_gateway from ..models.db import get_db from ..sim.ledger import Ledger from ..sim.matcher import Matcher -from .clock import can_open_new, window_key +from .clock import window_key from .exits import check_exits from .group import next_group_id -from .signal import decide logger = logging.getLogger(__name__) @@ -33,15 +32,18 @@ class StrategyEngine: assert row is not None upl = self.matcher.unrealized() s = get_settings() - exit_pts = self.ledger.get_setting_float("exit_move_points", s.exit_move_points) + exit_pct = self.ledger.get_setting_float("exit_move_pct", s.exit_move_pct) rest_sec = self.ledger.get_setting_int("rest_seconds", s.rest_seconds) - max_rounds = self.ledger.get_setting_int("max_rounds", s.max_rounds) + leverage = self.ledger.get_setting_float("leverage", s.leverage) + min_hours = self.ledger.get_setting_float("min_option_hours", s.min_option_hours) + min_opt_lev = self.ledger.get_setting_float( + "min_option_leverage", s.min_option_leverage + ) rest_until = row["rest_until_ms"] rest_left = 0 if rest_until: rest_left = max(0, int((int(rest_until) - time.time() * 1000) / 1000)) last_error = row["last_error"] - # 清掉已修复的旧序列化错误残留 if last_error and "PriceResult" in str(last_error) and "__dict__" in str(last_error): self._set_state(last_error=None) last_error = None @@ -49,13 +51,15 @@ class StrategyEngine: "running": bool(row["running"]), "phase": row["phase"], "rounds_done": int(row["rounds_done"] or 0), - "max_rounds": max_rounds, "window_key": row["window_key"], "rest_until_ms": rest_until, "rest_left_sec": rest_left, "rest_seconds": rest_sec, - "exit_move_points": exit_pts, - "can_open": can_open_new(open_hhmm=s.open_hhmm, stop_hhmm=s.stop_open_hhmm), + "exit_move_pct": exit_pct, + "leverage": leverage, + "min_option_hours": min_hours, + "min_option_leverage": min_opt_lev, + "can_open": True, "last_error": last_error, "position": upl, "ledger": self.ledger.snapshot(), @@ -103,23 +107,14 @@ class StrategyEngine: assert row is not None rounds = int(row["rounds_done"] or 0) + 1 rest_sec = self.ledger.get_setting_int("rest_seconds", s.rest_seconds) - max_rounds = self.ledger.get_setting_int("max_rounds", s.max_rounds) rest_until = int(time.time() * 1000) + rest_sec * 1000 - if rounds >= max_rounds: - self._set_state( - rounds_done=rounds, - phase="stopped", - rest_until_ms=None, - ) - else: - self._set_state( - rounds_done=rounds, - phase="resting", - rest_until_ms=rest_until, - ) + self._set_state( + rounds_done=rounds, + phase="resting", + rest_until_ms=rest_until, + ) - def _count_groups_for_window(self, wkey: str) -> int: - # group_id like G-20260724-01 ; window_key is YYYYMMDD + def _count_groups_for_day(self, wkey: str) -> int: rows = self.db.fetchall( "SELECT group_id FROM groups WHERE group_id LIKE ?", (f"G-{wkey}-%",), @@ -135,12 +130,11 @@ class StrategyEngine: await asyncio.sleep(1) continue async with self._lock: - # 空仓且 ATM 偏离现价时先重选,再跑开仓逻辑 try: await get_gateway().ensure_atm_async(force=False) except Exception as e: logger.warning("ATM ensure before tick failed: %s", e) - await asyncio.to_thread(self._tick) + await self._tick_async() except asyncio.CancelledError: raise except Exception as e: @@ -148,34 +142,33 @@ class StrategyEngine: self._set_state(last_error=str(e)) await asyncio.sleep(1) - def _tick(self) -> None: + async def _tick_async(self) -> None: s = get_settings() st = self.db.fetchone("SELECT * FROM strategy_state WHERE id=1") assert st is not None wkey = window_key() if st["window_key"] != wkey: - # 新业务窗重置轮次 - self._set_state(window_key=wkey, rounds_done=0, phase="idle", rest_until_ms=None) + self._set_state(window_key=wkey, phase="idle") st = self.db.fetchone("SELECT * FROM strategy_state WHERE id=1") assert st is not None - max_rounds = self.ledger.get_setting_int("max_rounds", s.max_rounds) - exit_pts = self.ledger.get_setting_float("exit_move_points", s.exit_move_points) + exit_pct = self.ledger.get_setting_float("exit_move_pct", s.exit_move_pct) pos = self.matcher.current_position() - # 有仓:盯平仓 if pos.get("status") == "open": self._set_state(phase="open", last_error=None) upl = self.matcher.unrealized() decision = check_exits( perp_upl=float(upl["perp_upl"]), initial_premium=float(upl["initial_premium"] or 0), - move_points=float(upl["move_points"] or 0), - exit_move_points=exit_pts, + move_pct=float(upl.get("move_pct") or 0), + exit_move_pct=exit_pct, ) if decision.should_close: self._set_state(phase="closing") - r = self.matcher.close_group(reason=decision.reason) + r = await asyncio.to_thread( + self.matcher.close_group, reason=decision.reason + ) if r.ok: self._after_close() elif r.liquidity_wait: @@ -184,7 +177,6 @@ class StrategyEngine: self._set_state(last_error=r.detail) return - # 休息中 if st["phase"] == "resting" and st["rest_until_ms"]: if int(time.time() * 1000) < int(st["rest_until_ms"]): return @@ -192,46 +184,38 @@ class StrategyEngine: st = self.db.fetchone("SELECT * FROM strategy_state WHERE id=1") assert st is not None - if int(st["rounds_done"] or 0) >= max_rounds: - self._set_state(phase="stopped") + if st["phase"] in ("paused",): return + # 旧「轮次停开」状态:自动恢复为空闲以便继续 + if st["phase"] in ("stopped", "outside_window"): + self._set_state(phase="idle") - if not can_open_new(open_hhmm=s.open_hhmm, stop_hhmm=s.stop_open_hhmm): - self._set_state(phase="outside_window") - return - - if st["phase"] in ("stopped", "paused"): - return - - # 尝试开仓 self._set_state(phase="wait_signal") gw = get_gateway() - snap = gw.snapshot() - if not snap.pair or not snap.call or not snap.put: - return - sig = decide(snap.call.ask, snap.put.ask) - if sig is None: + pick = await gw.pick_for_open_async() + if pick is None: + self._set_state( + last_error="无合格期权:需剩余时长与杠杆倍数同时满足" + ) return - self._set_state(phase="opening") - count = self._count_groups_for_window(wkey) + self._set_state(phase="opening", last_error=None) + count = self._count_groups_for_day(wkey) gid = next_group_id(count) option_inst = ( - snap.pair.call_inst_id if sig.option_side == "call" else snap.pair.put_inst_id + pick.pair.call_inst_id if pick.option_side == "call" else pick.pair.put_inst_id ) - entry_idx = snap.index_px or (snap.perp.mark_px if snap.perp else None) - if entry_idx is None: - self._set_state(last_error="no index/mark for entry") - return - r = self.matcher.open_group( + entry_idx = pick.underlying_px + r = await asyncio.to_thread( + self.matcher.open_group, group_id=gid, - bias=sig.bias, - option_side=sig.option_side, - perp_side=sig.perp_side, + bias=pick.bias, + option_side=pick.option_side, + perp_side=pick.perp_side, option_inst_id=option_inst, entry_index_px=float(entry_idx), - strike=snap.pair.strike, - expiry_ymd=snap.pair.expiry_ymd, + strike=pick.pair.strike, + expiry_ymd=pick.pair.expiry_ymd, ) if r.ok: self._set_state(phase="open", last_error=None) diff --git a/backend/app/strategy/exits.py b/backend/app/strategy/exits.py index 68b216c..5c0b746 100644 --- a/backend/app/strategy/exits.py +++ b/backend/app/strategy/exits.py @@ -13,11 +13,11 @@ def check_exits( *, perp_upl: float, initial_premium: float, - move_points: float, - exit_move_points: float, + move_pct: float, + exit_move_pct: float, ) -> ExitDecision: if initial_premium > 0 and perp_upl + 1e-9 >= initial_premium: return ExitDecision(True, "premium_cover") - if exit_move_points > 0 and move_points + 1e-9 >= exit_move_points: - return ExitDecision(True, "move_points") + if exit_move_pct > 0 and move_pct + 1e-9 >= exit_move_pct: + return ExitDecision(True, "move_pct") return ExitDecision(False, "") diff --git a/backend/tests/test_instruments.py b/backend/tests/test_instruments.py index 84bf880..492344a 100644 --- a/backend/tests/test_instruments.py +++ b/backend/tests/test_instruments.py @@ -4,7 +4,10 @@ from datetime import datetime from zoneinfo import ZoneInfo from app.market.instruments import ( + hours_until_expiry, + list_eligible_expiry_ymds, next_session_expiry_ymd, + option_leverage, parse_option_inst_id, pick_atm_strike, select_option_pair, @@ -48,3 +51,26 @@ def test_select_option_pair_atm() -> None: assert pair.strike == 3500 assert pair.call_inst_id.endswith("-3500-C") assert pair.put_inst_id.endswith("-3500-P") + + +def test_eligible_skips_short_ttm() -> None: + # 2026-07-25 08:30 SH:当日 16:00 到期仅约 7.5h,应跳过 260725,选 260726 + now = datetime(2026, 7, 25, 8, 30, tzinfo=_SH) + rows = [ + {"instId": "ETH-USD_UM-260725-1850-C", "state": "live"}, + {"instId": "ETH-USD_UM-260725-1850-P", "state": "live"}, + {"instId": "ETH-USD_UM-260726-1850-C", "state": "live"}, + {"instId": "ETH-USD_UM-260726-1850-P", "state": "live"}, + ] + assert hours_until_expiry("260725", now) < 12 + assert hours_until_expiry("260726", now) >= 12 + eligible = list_eligible_expiry_ymds(rows, min_hours=12, now=now) + assert eligible == ["260726"] + pair = select_option_pair(rows, mark_px=1853, min_hours=12, now=now) + assert pair is not None + assert pair.expiry_ymd == "260726" + + +def test_option_leverage() -> None: + assert abs((option_leverage(1850, 18.5) or 0) - 100) < 1e-9 + assert option_leverage(1850, 0) is None diff --git a/backend/tests/test_p1_p2_rules.py b/backend/tests/test_p1_p2_rules.py index 4abb3de..4bf8e0d 100644 --- a/backend/tests/test_p1_p2_rules.py +++ b/backend/tests/test_p1_p2_rules.py @@ -1,10 +1,11 @@ -from app.strategy.signal import decide -from app.strategy.exits import check_exits -from app.sim.pricing import option_fill, perp_fill -from app.strategy.clock import can_open_new, window_key from datetime import datetime from zoneinfo import ZoneInfo +from app.sim.pricing import option_fill, perp_fill +from app.strategy.clock import can_open_new, window_key +from app.strategy.exits import check_exits +from app.strategy.signal import decide + _SH = ZoneInfo("Asia/Shanghai") @@ -26,13 +27,19 @@ def test_signal_equal() -> None: assert decide(10.0, 10.0) is None -def test_exit_premium_and_move() -> None: - assert check_exits( - perp_upl=50, initial_premium=40, move_points=1, exit_move_points=30 - ).reason == "premium_cover" - assert check_exits( - perp_upl=1, initial_premium=40, move_points=30, exit_move_points=30 - ).reason == "move_points" +def test_exit_premium_and_move_pct() -> None: + assert ( + check_exits( + perp_upl=50, initial_premium=40, move_pct=0.1, exit_move_pct=2 + ).reason + == "premium_cover" + ) + assert ( + check_exits( + perp_upl=1, initial_premium=40, move_pct=2.0, exit_move_pct=2 + ).reason + == "move_pct" + ) def test_perp_pricing() -> None: @@ -47,15 +54,10 @@ def test_option_open_close_pricing() -> None: assert c.fill_px < 10 -def test_window() -> None: - # 17:00 can open, window key today +def test_window_always_open() -> None: n = datetime(2026, 7, 24, 17, 0, tzinfo=_SH) assert can_open_new(n) is True assert window_key(n) == "20260724" - # 10:00 cannot open n2 = datetime(2026, 7, 24, 10, 0, tzinfo=_SH) - assert can_open_new(n2) is False - # 07:00 still previous window, can open - n3 = datetime(2026, 7, 24, 7, 0, tzinfo=_SH) - assert can_open_new(n3) is True - assert window_key(n3) == "20260723" + assert can_open_new(n2) is True + assert window_key(n2) == "20260724" diff --git a/frontend/src/api/client.ts b/frontend/src/api/client.ts index 01a5daa..8b0f23d 100644 --- a/frontend/src/api/client.ts +++ b/frontend/src/api/client.ts @@ -117,11 +117,13 @@ export type PlanState = { running: boolean; phase: string; rounds_done: number; - max_rounds: number; window_key: string | null; rest_left_sec: number; rest_seconds: number; - exit_move_points: number; + exit_move_pct: number; + leverage: number; + min_option_hours: number; + min_option_leverage: number; can_open: boolean; last_error: string | null; position: { @@ -133,6 +135,9 @@ export type PlanState = { perp_entry_px?: number; perp_qty_eth?: number; perp_mark_px?: number | null; + perp_notional?: number; + perp_margin?: number | null; + leverage?: number; option_inst_id?: string; option_entry_px?: number; option_qty_eth?: number; @@ -145,6 +150,7 @@ export type PlanState = { index_px?: number | null; entry_index_px?: number; move_points?: number; + move_pct?: number; initial_premium?: number; premium_gap?: number; }; @@ -153,10 +159,12 @@ export type PlanState = { export type StrategySettings = { fee_rate: number; - exit_move_points: number; + exit_move_pct: number; rest_seconds: number; - max_rounds: number; initial_equity: number; + leverage: number; + min_option_hours: number; + min_option_leverage: number; perp_qty_eth: number; option_qty_eth: number; ledger: { equity: number; available: number }; diff --git a/frontend/src/pages/Plan.tsx b/frontend/src/pages/Plan.tsx index 9f3d2ad..0cf978e 100644 --- a/frontend/src/pages/Plan.tsx +++ b/frontend/src/pages/Plan.tsx @@ -75,15 +75,15 @@ export default function PlanPage() { const pos = plan?.position; const open = !!pos?.has_position; - const exitN = plan?.exit_move_points ?? 30; - const move = pos?.move_points ?? 0; + const exitPct = plan?.exit_move_pct ?? 2; + const movePct = pos?.move_pct ?? 0; const phaseLabel = PHASE_ZH[plan?.phase || ""] || plan?.phase || "—"; return (

自动对冲计划

- SIM 本地撮合 · 期权只买 · 标的波动 N 点可配 + SIM 本地撮合 · 期权只买 · 剩余时长选到期 · 波动按比例全平

{err ?
{err}
: null} @@ -142,14 +142,14 @@ export default function PlanPage() { )} {" "} - · {phaseLabel} · 轮次 {plan?.rounds_done ?? 0}/{plan?.max_rounds ?? 3} + · {phaseLabel} · 已完成 {plan?.rounds_done ?? 0} 组
- 开仓窗 + 选约条件 - {plan?.can_open ? "可开" : "禁止新开"} · 窗 {plan?.window_key || "—"} + 剩余≥{fmt(plan?.min_option_hours, 0)}h · 期权杠杆≥{fmt(plan?.min_option_leverage, 0)}x
@@ -159,9 +159,10 @@ export default function PlanPage() {
- 权益 + 权益 / 杠杆 - {fmt(plan?.ledger?.equity)} / 可用 {fmt(plan?.ledger?.available)} + {fmt(plan?.ledger?.equity)} / 可用 {fmt(plan?.ledger?.available)} ·{" "} + {fmt(plan?.leverage, 0)}x
@@ -181,9 +182,9 @@ export default function PlanPage() {
- N 点进度 + 波动进度 - {fmt(move, 1)} / {fmt(exitN, 0)} + {fmt(movePct, 2)}% / {fmt(exitPct, 2)}%
{plan?.last_error ? ( @@ -217,6 +218,9 @@ export default function PlanPage() { 永续持仓 组 {pos?.group_id} 数量 {fmt(pos?.perp_qty_eth, 2)} ETH + + {fmt(pos?.leverage ?? plan?.leverage, 0)}x +
@@ -234,16 +238,16 @@ export default function PlanPage() {
- 指数 - {fmt(pos?.index_px)} + 保证金 + {fmt(pos?.perp_margin)}
- 开仓指数 - {fmt(pos?.entry_index_px)} + 名义价值 + {fmt(pos?.perp_notional)}
- 波动点数 - {fmt(move, 1)} + 波动比例 + {fmt(movePct, 2)}%
diff --git a/frontend/src/pages/Settings.tsx b/frontend/src/pages/Settings.tsx index 7209fb7..49e2706 100644 --- a/frontend/src/pages/Settings.tsx +++ b/frontend/src/pages/Settings.tsx @@ -21,9 +21,11 @@ export default function SettingsPage() { const [loading, setLoading] = useState(false); const [fee, setFee] = useState(0.0005); - const [exitPts, setExitPts] = useState(30); + const [exitPct, setExitPct] = useState(2); const [rest, setRest] = useState(300); - const [maxRounds, setMaxRounds] = useState(3); + const [leverage, setLeverage] = useState(3); + const [minHours, setMinHours] = useState(12); + const [minOptLev, setMinOptLev] = useState(100); const [perpQty, setPerpQty] = useState(1); const [optQty, setOptQty] = useState(2); const [stratOk, setStratOk] = useState(""); @@ -32,9 +34,11 @@ export default function SettingsPage() { apiFetch("/api/settings/strategy") .then((s) => { setFee(s.fee_rate); - setExitPts(s.exit_move_points); + setExitPct(s.exit_move_pct ?? 2); setRest(s.rest_seconds); - setMaxRounds(s.max_rounds); + setLeverage(s.leverage ?? 3); + setMinHours(s.min_option_hours ?? 12); + setMinOptLev(s.min_option_leverage ?? 100); setPerpQty(s.perp_qty_eth ?? 1); setOptQty(s.option_qty_eth ?? 2); }) @@ -81,9 +85,11 @@ export default function SettingsPage() { method: "PUT", body: JSON.stringify({ fee_rate: fee, - exit_move_points: exitPts, + exit_move_pct: exitPct, rest_seconds: rest, - max_rounds: maxRounds, + leverage, + min_option_hours: minHours, + min_option_leverage: minOptLev, perp_qty_eth: perpQty, option_qty_eth: optQty, }), @@ -117,11 +123,62 @@ export default function SettingsPage() { {tab === "strategy" ? (

- 标的波动 N 点全平、轮次休息、仓位数量、费率(滑点=1×费率)。 + 无开仓时间窗;期权按剩余时长选到期 → 平值 → 校验杠杆;波动按百分比全平。

{stratOk ?
{stratOk}
: null} {err && tab === "strategy" ?
{err}
: null}
+

资金

+
+ + setLeverage(Number(e.target.value))} + /> +
+ +

策略

+
+ + setMinHours(Number(e.target.value))} + /> +
+
+ + setMinOptLev(Number(e.target.value))} + /> +
+
+ + setExitPct(Number(e.target.value))} + /> +
setOptQty(Number(e.target.value))} />
-
- - setExitPts(Number(e.target.value))} - /> -
setRest(Number(e.target.value))} />
-
- - setMaxRounds(Number(e.target.value))} - /> -