diff --git a/backend_api_python/app/services/live_trading/account_configuration.py b/backend_api_python/app/services/live_trading/account_configuration.py index 50a58ee88..fe090179f 100644 --- a/backend_api_python/app/services/live_trading/account_configuration.py +++ b/backend_api_python/app/services/live_trading/account_configuration.py @@ -2,10 +2,13 @@ from __future__ import annotations +import logging from typing import Any, Dict from app.services.live_trading.base import LiveTradingError +logger = logging.getLogger(__name__) + def requires_derivatives_account_configuration(*, market_type: str, reduce_only: bool) -> bool: """Only opening derivative orders may mutate symbol account settings.""" @@ -136,15 +139,31 @@ def configure_derivatives_account( details["account_mode"] = account_level if account_level == "1": raise LiveTradingError("OKX_SWAP_ACCOUNT_MODE_REQUIRED") + + inst_id = to_okx_swap_inst_id(symbol) + + # Proactively cancel algo orders that would block leverage changes. + # This prevents the common 59669 error on strategy restart where stale + # algo orders from a previous run remain on the exchange. + try: + cancelled = client.cancel_all_algo_orders(inst_id=inst_id, inst_type="SWAP") + if cancelled > 0: + logger.info( + f"Pre-cleaned {cancelled} lingering algo order(s) on {inst_id} " + "before setting leverage." + ) + except Exception: + pass + if position_mode in ("long_short_mode", "longshort_mode"): long_ok = client.set_leverage( - inst_id=to_okx_swap_inst_id(symbol), + inst_id=inst_id, lever=target_leverage, mgn_mode=mode, pos_side="long", ) short_ok = client.set_leverage( - inst_id=to_okx_swap_inst_id(symbol), + inst_id=inst_id, lever=target_leverage, mgn_mode=mode, pos_side="short", @@ -153,7 +172,7 @@ def configure_derivatives_account( details["position_mode"] = "hedge" else: ok = client.set_leverage( - inst_id=to_okx_swap_inst_id(symbol), + inst_id=inst_id, lever=target_leverage, mgn_mode=mode, pos_side="net", diff --git a/backend_api_python/app/services/live_trading/external_flat_close.py b/backend_api_python/app/services/live_trading/external_flat_close.py new file mode 100644 index 000000000..387281a68 --- /dev/null +++ b/backend_api_python/app/services/live_trading/external_flat_close.py @@ -0,0 +1,473 @@ +"""Record strategy close trades when exchange flat-closes a local leg. + +Native TP/SL (algo) / liquidation / manual exchange closes often never create a +``pending_orders`` row, so the private-stream projector cannot attribute them. +Position sync already purges the ghost L3 row; this module also writes the +matching ``close_*`` trade so win-rate / trade history stay honest. +""" + +from __future__ import annotations + +from datetime import datetime, timezone +from typing import Any, Dict, Optional, Tuple + +from app.services.live_trading.okx import OkxClient +from app.services.live_trading.records import ( + _fetch_position, + apply_fill_to_local_position, + normalize_strategy_symbol, + record_trade, +) +from app.services.live_trading.symbols import to_okx_swap_inst_id +from app.utils.db import get_db_connection +from app.utils.logger import get_logger +from app.utils.strategy_runtime_logs import append_strategy_log +from app.utils.trade_close_reason import ( + EXCHANGE_ADL, + EXCHANGE_FLAT_RECONCILE, + EXCHANGE_LIQUIDATION, + EXCHANGE_NATIVE_CLOSE, +) + +logger = get_logger(__name__) + +_EPS = 1e-12 + + +def trade_side_net_qty(strategy_id: int, symbol: str, side: str) -> float: + """Net base qty still open on the trade ledger for one side.""" + sid = int(strategy_id or 0) + sym = normalize_strategy_symbol(symbol) or str(symbol or "").strip() + side_l = str(side or "").strip().lower() + if sid <= 0 or not sym or side_l not in ("long", "short"): + return 0.0 + open_types = ("open_long", "add_long") if side_l == "long" else ("open_short", "add_short") + close_types = ("close_long", "reduce_long") if side_l == "long" else ("close_short", "reduce_short") + with get_db_connection() as db: + cur = db.cursor() + cur.execute( + """ + SELECT type, amount + FROM qd_strategy_trades + WHERE strategy_id = %s + AND UPPER(COALESCE(NULLIF(symbol_canonical, ''), symbol)) = UPPER(%s) + AND type = ANY(%s) + ORDER BY id ASC + """, + (sid, sym, list(open_types) + list(close_types)), + ) + rows = cur.fetchall() or [] + cur.close() + net = 0.0 + for row in rows: + t = str(row.get("type") or "").strip().lower() + try: + qty = float(row.get("amount") or 0.0) + except Exception: + qty = 0.0 + if t in open_types: + net += qty + elif t in close_types: + net -= qty + return net if net > _EPS else 0.0 + + +def _last_open_trade(strategy_id: int, symbol: str, side: str) -> Dict[str, Any]: + sid = int(strategy_id or 0) + sym = normalize_strategy_symbol(symbol) or str(symbol or "").strip() + side_l = str(side or "").strip().lower() + open_types = ("open_long", "add_long") if side_l == "long" else ("open_short", "add_short") + with get_db_connection() as db: + cur = db.cursor() + cur.execute( + """ + SELECT id, price, amount, created_at, credential_id, inst_id, market_type + FROM qd_strategy_trades + WHERE strategy_id = %s + AND UPPER(COALESCE(NULLIF(symbol_canonical, ''), symbol)) = UPPER(%s) + AND type = ANY(%s) + ORDER BY id DESC + LIMIT 1 + """, + (sid, sym, list(open_types)), + ) + row = cur.fetchone() or {} + cur.close() + return dict(row) if isinstance(row, dict) else {} + + +def _already_recorded_fill(strategy_id: int, exchange_fill_id: str) -> bool: + fill_id = str(exchange_fill_id or "").strip() + if not fill_id: + return False + with get_db_connection() as db: + cur = db.cursor() + cur.execute( + """ + SELECT 1 FROM qd_strategy_trades + WHERE strategy_id = %s AND exchange_fill_id = %s + LIMIT 1 + """, + (int(strategy_id), fill_id), + ) + exists = cur.fetchone() is not None + cur.close() + return bool(exists) + + +def _as_float(value: Any, default: float = 0.0) -> float: + try: + return float(value) + except Exception: + return float(default) + + +def _parse_okx_ms(value: Any) -> Optional[datetime]: + try: + ms = int(float(value)) + except Exception: + return None + if ms <= 0: + return None + return datetime.fromtimestamp(ms / 1000.0, tz=timezone.utc) + + +def resolve_okx_external_close( + client: Any, + *, + symbol: str, + side: str, + entry_price: float, + opened_after: Optional[datetime] = None, +) -> Optional[Dict[str, Any]]: + """Best-effort close snapshot from OKX positions-history / fills-history.""" + if not isinstance(client, OkxClient): + return None + sym = normalize_strategy_symbol(symbol) or str(symbol or "").strip() + side_l = str(side or "").strip().lower() + if not sym or side_l not in ("long", "short"): + return None + inst_id = to_okx_swap_inst_id(sym) + entry = _as_float(entry_price) + opened_ts = None + if isinstance(opened_after, datetime): + opened_ts = opened_after if opened_after.tzinfo else opened_after.replace(tzinfo=timezone.utc) + + # Prefer positions-history: one row per fully closed cycle with open/close px. + try: + hist = client._signed_request( + "GET", + "/api/v5/account/positions-history", + params={"instType": "SWAP", "instId": inst_id, "limit": "30"}, + ) + rows = (hist or {}).get("data") if isinstance(hist, dict) else [] + except Exception as exc: + logger.debug("OKX positions-history lookup failed %s: %s", inst_id, exc) + rows = [] + + best: Optional[Dict[str, Any]] = None + for row in rows or []: + if not isinstance(row, dict): + continue + pos_side = str(row.get("posSide") or "").strip().lower() + direction = str(row.get("direction") or "").strip().lower() + row_side = pos_side if pos_side in ("long", "short") else direction + if row_side and row_side != side_l: + continue + open_px = _as_float(row.get("openAvgPx")) + close_px = _as_float(row.get("closeAvgPx")) + if close_px <= 0: + continue + if entry > 0 and open_px > 0: + # Match the cycle that started at our local entry. + if abs(open_px - entry) / max(entry, 1.0) > 0.0005 and abs(open_px - entry) > 1.0: + continue + closed_at = _parse_okx_ms(row.get("uTime") or row.get("cTime")) + if opened_ts and closed_at and closed_at < opened_ts: + continue + close_type = str(row.get("type") or "").strip() + reason = EXCHANGE_NATIVE_CLOSE + # OKX type: 1 partial close, 2 full close, 3 liquidation, 4 partial liq, 5 ADL + if close_type in {"3", "4"}: + reason = EXCHANGE_LIQUIDATION + elif close_type == "5": + reason = EXCHANGE_ADL + candidate = { + "price": close_px, + "fee": 0.0, + "fee_ccy": "USDT", + "exchange_fill_id": "", + "closed_at": closed_at, + "close_reason": reason, + "source": "positions_history", + } + if best is None: + best = candidate + continue + prev_closed = best.get("closed_at") + if closed_at and (not prev_closed or closed_at > prev_closed): + best = candidate + + try: + fills = client._signed_request( + "GET", + "/api/v5/trade/fills-history", + params={"instType": "SWAP", "instId": inst_id, "limit": "50"}, + ) + fill_rows = (fills or {}).get("data") if isinstance(fills, dict) else [] + except Exception: + fill_rows = [] + + close_side = "sell" if side_l == "long" else "buy" + + # Enrich fee / fill id from recent fills when positions-history matched. + if best is not None: + target_px = _as_float(best.get("price")) + for fill in fill_rows or []: + if not isinstance(fill, dict): + continue + if str(fill.get("side") or "").strip().lower() != close_side: + continue + pos_side = str(fill.get("posSide") or "").strip().lower() + if pos_side and pos_side not in ("net", side_l): + continue + fill_px = _as_float(fill.get("fillPx")) + if target_px > 0 and fill_px > 0 and abs(fill_px - target_px) / target_px > 0.0005: + continue + fill_ts = _parse_okx_ms(fill.get("ts")) + if opened_ts and fill_ts and fill_ts < opened_ts: + continue + fee = abs(_as_float(fill.get("fee"))) + best["fee"] = fee + best["fee_ccy"] = str(fill.get("feeCcy") or "USDT") + best["exchange_fill_id"] = str(fill.get("tradeId") or fill.get("fillId") or "") + if fill_ts: + best["closed_at"] = fill_ts + break + return best + + # Race fallback: positions-history can lag a just-closed fill by a few seconds. + # Prefer the newest reduce-side fill after our open instead of inventing entry==close. + fill_best: Optional[Dict[str, Any]] = None + for fill in fill_rows or []: + if not isinstance(fill, dict): + continue + if str(fill.get("side") or "").strip().lower() != close_side: + continue + pos_side = str(fill.get("posSide") or "").strip().lower() + if pos_side and pos_side not in ("net", side_l): + continue + fill_px = _as_float(fill.get("fillPx")) + if fill_px <= 0: + continue + fill_ts = _parse_okx_ms(fill.get("ts")) + if opened_ts and fill_ts and fill_ts < opened_ts: + continue + candidate = { + "price": fill_px, + "fee": abs(_as_float(fill.get("fee"))), + "fee_ccy": str(fill.get("feeCcy") or "USDT"), + "exchange_fill_id": str(fill.get("tradeId") or fill.get("fillId") or ""), + "closed_at": fill_ts, + "close_reason": EXCHANGE_NATIVE_CLOSE, + "source": "fills_history", + } + if fill_best is None: + fill_best = candidate + continue + prev_closed = fill_best.get("closed_at") + if fill_ts and (not prev_closed or fill_ts > prev_closed): + fill_best = candidate + return fill_best + + +def _gross_close_profit(*, side: str, entry_price: float, close_price: float, amount: float) -> float: + entry = _as_float(entry_price) + close = _as_float(close_price) + qty = _as_float(amount) + if entry <= 0 or close <= 0 or qty <= 0: + return 0.0 + if str(side).strip().lower() == "long": + return (close - entry) * qty + return (entry - close) * qty + + +def record_external_flat_close( + *, + strategy_id: int, + symbol: str, + side: str, + amount: float, + entry_price: float, + close_price: float, + commission: float = 0.0, + commission_ccy: str = "USDT", + close_reason: str = EXCHANGE_NATIVE_CLOSE, + fill_source: str = "position_sync", + exchange_fill_id: str = "", + credential_id: int = 0, + inst_id: str = "", + market_type: str = "swap", + created_at: Optional[datetime] = None, + clear_local_position: bool = True, +) -> int: + """Persist one external close trade; optionally clear the local L3 leg.""" + sid = int(strategy_id or 0) + sym = normalize_strategy_symbol(symbol) or str(symbol or "").strip() + side_l = str(side or "").strip().lower() + qty = _as_float(amount) + px = _as_float(close_price) + if sid <= 0 or not sym or side_l not in ("long", "short") or qty <= _EPS or px <= 0: + return 0 + if exchange_fill_id and _already_recorded_fill(sid, exchange_fill_id): + return 0 + + trade_type = "close_long" if side_l == "long" else "close_short" + reason = str(close_reason or EXCHANGE_FLAT_RECONCILE).strip() or EXCHANGE_FLAT_RECONCILE + entry = _as_float(entry_price) + profit = _gross_close_profit(side=side_l, entry_price=entry, close_price=px, amount=qty) + fee = abs(_as_float(commission)) + + trade_id = record_trade( + strategy_id=sid, + symbol=sym, + trade_type=trade_type, + price=px, + amount=qty, + commission=fee, + commission_ccy=str(commission_ccy or "USDT"), + commission_quote=fee if str(commission_ccy or "USDT").upper() in ("", "USDT", "USD") else None, + profit=profit, + close_reason=reason, + matched_entry_price=entry if entry > 0 else None, + fill_source=str(fill_source or "position_sync"), + exchange_fill_id=str(exchange_fill_id or ""), + credential_id=int(credential_id or 0), + inst_id=str(inst_id or ""), + market_type=str(market_type or "swap"), + fee_status="actual" if fee > 0 else "pending", + fee_source="rest" if fee > 0 else "", + created_at=created_at, + ) + if trade_id > 0 and clear_local_position: + local = _fetch_position(sid, sym, side_l) + local_size = _as_float(local.get("size")) + if local_size > _EPS: + apply_fill_to_local_position( + strategy_id=sid, + symbol=sym, + signal_type=trade_type, + filled=min(local_size, qty), + avg_price=px, + ) + if trade_id > 0: + try: + append_strategy_log( + sid, + "trade", + ( + f"Trade reconciled: {trade_type} {sym} filled={qty:.8f} @ {px:.8f}, " + f"profit={profit:.4f}, reason={reason} (source={fill_source})" + ), + ) + except Exception: + pass + return int(trade_id or 0) + + +def reconcile_external_flat_closes( + *, + strategy_id: int, + symbol: str, + side: str, + client: Any = None, + fallback_price: float = 0.0, + credential_id: int = 0, + inst_id: str = "", + market_type: str = "swap", + clear_local_position: bool = True, +) -> int: + """ + If the trade ledger still shows open qty for a leg that is flat on the + exchange, write the missing close and optionally clear local size. + """ + sid = int(strategy_id or 0) + sym = normalize_strategy_symbol(symbol) or str(symbol or "").strip() + side_l = str(side or "").strip().lower() + residual = trade_side_net_qty(sid, sym, side_l) + if residual <= _EPS: + local = _fetch_position(sid, sym, side_l) + local_size = _as_float(local.get("size")) + if local_size <= _EPS: + return 0 + residual = local_size + entry = _as_float(local.get("entry_price")) + cred = int(local.get("credential_id") or credential_id or 0) + iid = str(local.get("inst_id") or inst_id or "") + mt = str(local.get("market_type") or market_type or "swap") + opened_after = None + else: + last_open = _last_open_trade(sid, sym, side_l) + entry = _as_float(last_open.get("price")) + if entry <= 0: + local = _fetch_position(sid, sym, side_l) + entry = _as_float(local.get("entry_price")) + cred = int(last_open.get("credential_id") or credential_id or 0) + iid = str(last_open.get("inst_id") or inst_id or "") + mt = str(last_open.get("market_type") or market_type or "swap") + opened_after = last_open.get("created_at") + if isinstance(opened_after, datetime) and opened_after.tzinfo is None: + opened_after = opened_after.replace(tzinfo=timezone.utc) + + snapshot = resolve_okx_external_close( + client, + symbol=sym, + side=side_l, + entry_price=entry, + opened_after=opened_after if isinstance(opened_after, datetime) else None, + ) or {} + close_px = _as_float(snapshot.get("price")) + if close_px <= 0: + close_px = _as_float(fallback_price) + # Never invent close_price == entry_price: that creates fake 0-pnl rows with + # pending fees and confuses users. Skip and let a later sync retry. + if close_px <= 0: + logger.warning( + "[ExternalFlatClose] strategy=%s %s %s residual=%.8f but no close price available; deferring", + sid, + sym, + side_l, + residual, + ) + return 0 + if entry > 0 and abs(close_px - entry) <= 1e-9 and not snapshot: + logger.warning( + "[ExternalFlatClose] strategy=%s %s %s refusing entry-price fallback without exchange snapshot", + sid, + sym, + side_l, + ) + return 0 + + reason = str(snapshot.get("close_reason") or EXCHANGE_NATIVE_CLOSE) + if not snapshot: + reason = EXCHANGE_FLAT_RECONCILE + + return record_external_flat_close( + strategy_id=sid, + symbol=sym, + side=side_l, + amount=residual, + entry_price=entry, + close_price=close_px, + commission=_as_float(snapshot.get("fee")), + commission_ccy=str(snapshot.get("fee_ccy") or "USDT"), + close_reason=reason, + fill_source="position_sync", + exchange_fill_id=str(snapshot.get("exchange_fill_id") or ""), + credential_id=cred, + inst_id=iid, + market_type=mt, + created_at=snapshot.get("closed_at") if isinstance(snapshot.get("closed_at"), datetime) else None, + clear_local_position=clear_local_position, + ) diff --git a/backend_api_python/app/services/live_trading/okx.py b/backend_api_python/app/services/live_trading/okx.py index 8b1100f92..8ac13e819 100644 --- a/backend_api_python/app/services/live_trading/okx.py +++ b/backend_api_python/app/services/live_trading/okx.py @@ -13,7 +13,7 @@ import logging import time from decimal import Decimal, ROUND_DOWN -from typing import Any, Dict, Optional, Tuple +from typing import Any, Dict, List, Optional, Sequence, Tuple, Union from urllib.parse import urlencode from app.services.live_trading.base import BaseRestClient, LiveOrderResult, LiveTradingError @@ -290,7 +290,7 @@ def _signed_request( method: str, path: str, *, - json_body: Optional[Dict[str, Any]] = None, + json_body: Optional[Union[Dict[str, Any], List[Any]]] = None, params: Optional[Dict[str, Any]] = None, ) -> Dict[str, Any]: """ @@ -477,6 +477,88 @@ def get_positions(self, *, inst_id: str = "", inst_type: str = "SWAP") -> Dict[s params["instId"] = str(inst_id).strip() return self._signed_request("GET", "/api/v5/account/positions", params=params) + def _read_effective_leverage( + self, *, inst_id: str, mgn_mode: str = "cross", pos_side: str = "" + ) -> Optional[int]: + """Best-effort read of current leverage from an open position row.""" + iid = str(inst_id or "").strip() + if not iid: + return None + try: + resp = self.get_positions(inst_id=iid, inst_type="SWAP") + except Exception: + return None + rows = (resp.get("data") or []) if isinstance(resp, dict) else [] + mm = str(mgn_mode or "cross").strip().lower() + ps = str(pos_side or "").strip().lower() + for row in rows: + if not isinstance(row, dict): + continue + if str(row.get("instId") or "") != iid: + continue + if str(row.get("mgnMode") or "").strip().lower() != mm: + continue + row_ps = str(row.get("posSide") or "").strip().lower() + if ps and row_ps and row_ps not in (ps, "net"): + continue + try: + lev = int(float(row.get("lever") or 0)) + except (TypeError, ValueError): + continue + if lev > 0: + return lev + return None + + def get_leverage_info( + self, *, inst_id: str, mgn_mode: str = "cross" + ) -> Dict[str, Any]: + """Read configured leverage for an instrument (works without open position).""" + params: Dict[str, Any] = { + "instId": str(inst_id or "").strip(), + "mgnMode": str(mgn_mode or "cross").strip().lower(), + } + return self._signed_request("GET", "/api/v5/account/leverage-info", params=params) + + def _read_configured_leverage( + self, *, inst_id: str, mgn_mode: str = "cross", pos_side: str = "" + ) -> Optional[int]: + """Read leverage from open position or account leverage-info.""" + current = self._read_effective_leverage( + inst_id=inst_id, mgn_mode=mgn_mode, pos_side=pos_side + ) + if current is not None: + return current + iid = str(inst_id or "").strip() + if not iid: + return None + mm = str(mgn_mode or "cross").strip().lower() + ps = str(pos_side or "").strip().lower() + try: + resp = self.get_leverage_info(inst_id=iid, mgn_mode=mm) + except Exception: + return None + rows = (resp.get("data") or []) if isinstance(resp, dict) else [] + fallback: Optional[int] = None + for row in rows: + if not isinstance(row, dict): + continue + if str(row.get("instId") or "") != iid: + continue + try: + lev = int(float(row.get("lever") or 0)) + except (TypeError, ValueError): + continue + if lev <= 0: + continue + row_ps = str(row.get("posSide") or "").strip().lower() + if ps and row_ps and row_ps not in (ps, "net", ""): + continue + if ps and row_ps == ps: + return lev + if fallback is None: + fallback = lev + return fallback + def set_leverage(self, *, inst_id: str, lever: float, mgn_mode: str = "cross", pos_side: str = "") -> bool: """ Set leverage for an instrument (best-effort). @@ -487,6 +569,11 @@ def set_leverage(self, *, inst_id: str, lever: float, mgn_mode: str = "cross", p - lever - mgnMode: cross / isolated - posSide: net / long / short (required depending on posMode) + + On OKX error 59669 (algo orders blocking leverage change), automatically + cancels all outstanding algo orders and retries up to 3 times with + increasing backoff. After each cancellation, verifies that no algo orders + remain before retrying. """ iid = str(inst_id or "").strip() if not iid: @@ -503,8 +590,6 @@ def set_leverage(self, *, inst_id: str, lever: float, mgn_mode: str = "cross", p mm = "cross" ps = str(pos_side or "").strip().lower() - # In net_mode, OKX requires posSide=net. In long_short_mode, requires long/short. - # Caller should pass already resolved posSide; but keep a safe fallback. if ps not in ("net", "long", "short"): try: cfg = self.get_account_config() or {} @@ -521,12 +606,92 @@ def set_leverage(self, *, inst_id: str, lever: float, mgn_mode: str = "cross", p if ok and (now - float(ts or 0.0)) <= float(self._lev_cache_ttl_sec or 60.0): return True + current_lev = self._read_configured_leverage(inst_id=iid, mgn_mode=mm, pos_side=ps) + if current_lev == lv: + logger.debug( + "OKX leverage already %sx on %s (%s/%s); skipping set-leverage.", + lv, + iid, + mm, + ps or "net", + ) + self._lev_cache[cache_key] = (now, True) + return True + body: Dict[str, Any] = {"instId": iid, "lever": str(lv), "mgnMode": mm} if ps: body["posSide"] = ps - response = self._signed_request( - "POST", "/api/v5/account/set-leverage", json_body=body - ) + + MAX_RETRIES = 3 + for attempt in range(MAX_RETRIES + 1): + try: + response = self._signed_request( + "POST", "/api/v5/account/set-leverage", json_body=body + ) + break + except LiveTradingError as e: + error_text = str(e) + if "59669" in error_text or "Cancel cross-margin trailing" in error_text: + if attempt >= MAX_RETRIES: + if self._read_configured_leverage( + inst_id=iid, mgn_mode=mm, pos_side=ps + ) == lv: + logger.info( + "OKX 59669 on %s but configured leverage is already %sx; continuing.", + iid, + lv, + ) + self._lev_cache[cache_key] = (now, True) + response = {"data": [{"lever": str(lv)}]} + break + + error_detail = ( + f"OKX error 59669 persists after {MAX_RETRIES + 1} attempts " + f"(incl. retries) to cancel algo orders for {iid}. " + "Please manually check and cancel any trailing stop, trigger, " + "iceberg, TWAP, or chase orders on OKX for this instrument, " + "then restart the strategy." + ) + logger.error(error_detail) + raise LiveTradingError(error_detail) from e + + backoff = 1.0 + attempt + logger.warning( + f"OKX error 59669 (attempt {attempt + 1}/{MAX_RETRIES + 1}): " + f"algo orders blocking leverage change on {iid}. " + f"Cancelling all algo orders and retrying in {backoff}s..." + ) + + # Cancel algo orders on this specific instrument first + cancelled = self.cancel_all_algo_orders(inst_id=iid, inst_type="SWAP") + logger.info(f"Cancelled {cancelled} algo order(s) for {iid}.") + + # Also try cancelling ALL algo orders across all instruments, + # since OKX may block leverage changes for any instrument + if cancelled == 0 or attempt >= 1: + all_cancelled = self.cancel_all_algo_orders(inst_type="SWAP") + logger.info( + f"Cancelled {all_cancelled} algo order(s) across all instruments." + ) + cancelled = max(cancelled, all_cancelled) + + if cancelled == 0: + logger.warning( + f"No algo orders found to cancel for {iid}, " + "but OKX still reports 59669. This may indicate " + "stale algo orders or a different instrument blocking " + "the leverage change." + ) + + time.sleep(backoff) + else: + raise + else: + # This shouldn't happen in practice, but guard against it + raise LiveTradingError( + f"OKX set_leverage failed after {MAX_RETRIES + 1} attempts on {iid}" + ) + rows = (response.get("data") or []) if isinstance(response, dict) else [] first = rows[0] if isinstance(rows, list) and rows and isinstance(rows[0], dict) else {} effective_raw = first.get("lever") if isinstance(first, dict) else None @@ -541,6 +706,7 @@ def set_leverage(self, *, inst_id: str, lever: float, mgn_mode: str = "cross", p raise LiveTradingError( f"OKX applied {effective}x instead of requested {lv}x leverage" ) + self._lev_cache[cache_key] = (now, True) return True @@ -760,6 +926,151 @@ def cancel_order(self, *, market_type: str, symbol: str, ord_id: str = "", cl_or raise LiveTradingError("OKX cancel_order requires ord_id or cl_ord_id") return self._signed_request("POST", "/api/v5/trade/cancel-order", json_body=body) + def get_algo_orders(self, *, ord_type: str = "", inst_type: str = "SWAP", inst_id: str = "") -> Dict[str, Any]: + """ + Get outstanding algo orders (trigger, trailing stop, iceberg, TWAP, chase). + + Endpoint: GET /api/v5/trade/orders-algo-pending + Args: + ord_type: Filter by order type (conditional, oco, trigger, move_order_stop, + iceberg, twap, chase). Empty to return all types. + inst_type: SWAP / SPOT / FUTURES / OPTION (default: SWAP) + inst_id: Optional instrument ID filter (e.g. BTC-USDT-SWAP) + """ + params: Dict[str, Any] = {"ordType": ord_type, "instType": inst_type} + if inst_id: + params["instId"] = str(inst_id).strip() + resp = self._signed_request("GET", "/api/v5/trade/orders-algo-pending", params=params) + return resp if isinstance(resp, dict) else {"raw": resp} + + def cancel_algo_order(self, *, algo_id: str, inst_id: str, ord_type: str = "") -> Dict[str, Any]: + """ + Cancel an algo order (trigger, trailing stop, iceberg, TWAP, chase). + + OKX requires the body to be a JSON array. Trailing/iceberg/TWAP orders + use /trade/cancel-advance-algos; other algo types use /trade/cancel-algos. + """ + payload = [{"algoId": str(algo_id), "instId": str(inst_id)}] + return self._cancel_algo_payload(payload, ord_type=ord_type) + + def _cancel_algo_payload( + self, + payload: Sequence[Dict[str, str]], + *, + ord_type: str = "", + ) -> Dict[str, Any]: + items = [ + {"algoId": str(item.get("algoId") or ""), "instId": str(item.get("instId") or "")} + for item in payload + if str(item.get("algoId") or "") and str(item.get("instId") or "") + ] + if not items: + raise LiveTradingError("OKX cancel algo payload is empty") + + advanced_types = {"move_order_stop", "iceberg", "twap", "smart_iceberg"} + otype = str(ord_type or "").strip().lower() + endpoints: List[str] = [] + if otype in advanced_types: + endpoints = ["/api/v5/trade/cancel-advance-algos", "/api/v5/trade/cancel-algos"] + elif otype: + endpoints = ["/api/v5/trade/cancel-algos", "/api/v5/trade/cancel-advance-algos"] + else: + endpoints = ["/api/v5/trade/cancel-algos", "/api/v5/trade/cancel-advance-algos"] + + last_error: Optional[Exception] = None + for path in endpoints: + try: + return self._signed_request("POST", path, json_body=list(items)) + except Exception as exc: + last_error = exc + logger.debug("OKX %s failed for ordType=%s: %s", path, otype or "*", exc) + assert last_error is not None + raise last_error + + def cancel_all_algo_orders(self, *, inst_id: str = "", inst_type: str = "SWAP") -> int: + """ + Cancel all outstanding algo orders, optionally filtered by instrument. + + Args: + inst_id: Instrument ID filter (e.g. BTC-USDT-SWAP). Empty to cancel all. + inst_type: SWAP / SPOT / FUTURES / OPTION (default: SWAP) + + Returns: + Number of algo orders successfully cancelled. + """ + # Algo order types that block leverage adjustment per OKX error 59669 + algo_types = ["conditional", "oco", "trigger", "move_order_stop", "iceberg", "twap", "chase"] + + cancelled_count = 0 + all_algo_ids: list[Dict[str, str]] = [] + + # Fetch all pending algo orders across all types + for algo_type in algo_types: + try: + resp = self.get_algo_orders(ord_type=algo_type, inst_type=inst_type, inst_id=inst_id) + data = (resp.get("data") or []) if isinstance(resp, dict) else [] + for item in data: + if isinstance(item, dict): + all_algo_ids.append({ + "algoId": str(item.get("algoId") or ""), + "instId": str(item.get("instId") or ""), + "ordType": algo_type, + }) + except Exception as e: + logger.warning(f"Failed to fetch {algo_type} algo orders: {e}") + + if not all_algo_ids: + logger.info("No outstanding algo orders found.") + return 0 + + logger.info(f"Found {len(all_algo_ids)} outstanding algo order(s). Cancelling...") + + # Cancel in batches of 10 (OKX max). Group by ordType so the correct + # cancel endpoint can be selected for trailing/iceberg/TWAP orders. + by_type: Dict[str, List[Dict[str, str]]] = {} + for algo_info in all_algo_ids: + if not algo_info.get("algoId") or not algo_info.get("instId"): + logger.warning(f"Skipping algo order with missing algoId/instId: {algo_info}") + continue + by_type.setdefault(algo_info["ordType"], []).append( + {"algoId": algo_info["algoId"], "instId": algo_info["instId"]} + ) + + for otype, items in by_type.items(): + for idx in range(0, len(items), 10): + batch = items[idx : idx + 10] + try: + resp = self._cancel_algo_payload(batch, ord_type=otype) + data = (resp.get("data") or []) if isinstance(resp, dict) else [] + batch_ok = 0 + for row in data: + if not isinstance(row, dict): + continue + scode = str(row.get("sCode") or "") + if scode in ("0", ""): + batch_ok += 1 + else: + logger.warning( + "OKX refused algo cancel algoId=%s sCode=%s sMsg=%s", + row.get("algoId"), + scode, + row.get("sMsg"), + ) + if not data and str(resp.get("code") or "") == "0": + batch_ok = len(batch) + cancelled_count += batch_ok + logger.info( + "Cancelled %s/%s %s algo order(s) in batch", + batch_ok, + len(batch), + otype, + ) + except Exception as e: + logger.warning(f"Failed to cancel {otype} algo batch ({len(batch)}): {e}") + + logger.info(f"Successfully cancelled {cancelled_count}/{len(all_algo_ids)} algo order(s).") + return cancelled_count + def get_order(self, *, inst_id: str, ord_id: str = "", cl_ord_id: str = "") -> Dict[str, Any]: params: Dict[str, Any] = {"instId": str(inst_id)} if ord_id: diff --git a/backend_api_python/app/services/live_trading/records.py b/backend_api_python/app/services/live_trading/records.py index 0cfa54aeb..25cc80b1f 100644 --- a/backend_api_python/app/services/live_trading/records.py +++ b/backend_api_python/app/services/live_trading/records.py @@ -505,6 +505,7 @@ def record_trade( fee_source: str = "", fees_by_ccy: Optional[Dict[str, float]] = None, exchange_order_id: str = "", + created_at: Optional[Any] = None, ) -> int: value = float(amount or 0.0) * float(price or 0.0) if user_id is None: @@ -539,7 +540,7 @@ def record_trade( strategy_run_id, order_intent_id, execution_event_id, exchange_fill_id, fee_status, fee_source, commission_breakdown, exchange_order_id, created_at) VALUES - (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, NOW()) + (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, COALESCE(%s, NOW())) ON CONFLICT (execution_event_id) WHERE execution_event_id > 0 DO NOTHING RETURNING id """, @@ -573,6 +574,7 @@ def record_trade( str(fee_source or ""), json.dumps(fees_by_ccy or ({commission_ccy: commission} if commission_ccy and commission_ccy != 'MIXED' else {})), str(exchange_order_id or ''), + created_at, ), ) row = cur.fetchone() diff --git a/backend_api_python/app/services/pending_order_position_sync.py b/backend_api_python/app/services/pending_order_position_sync.py index 960e50b51..53b006b8e 100644 --- a/backend_api_python/app/services/pending_order_position_sync.py +++ b/backend_api_python/app/services/pending_order_position_sync.py @@ -23,7 +23,12 @@ from app.services.live_trading.gate import GateStockClient, GateUsdtFuturesClient from app.services.live_trading.leg_context import credential_id_from_exchange_config from app.services.live_trading.okx import OkxClient -from app.services.live_trading.records import normalize_strategy_symbol, strategy_allowed_symbols +from app.services.live_trading.records import ( + _delete_position, + lookup_exchange_side_qty, + normalize_strategy_symbol, + strategy_allowed_symbols, +) from app.services.pending_orders.position_sync_cache import ( exchange_sync_backoff_sec, get_position_sync_snapshot, @@ -49,6 +54,214 @@ _POSITION_SYNC_FD_BACKOFF_UNTIL = 0.0 +def _local_strategy_position_legs( + strategy_id: int, + allowed_symbols: set[str], +) -> List[tuple[str, str]]: + allowed = { + normalize_strategy_symbol(str(symbol or "")).upper() + for symbol in allowed_symbols + if normalize_strategy_symbol(str(symbol or "")) + } + if not allowed: + return [] + + with get_db_connection() as db: + cur = db.cursor() + cur.execute( + """ + SELECT symbol, side + FROM qd_strategy_positions + WHERE strategy_id = %s + AND COALESCE(size, 0) > 0 + AND side IN ('long', 'short') + """, + (int(strategy_id),), + ) + rows = cur.fetchall() or [] + cur.close() + + output: List[tuple[str, str]] = [] + seen: set[tuple[str, str]] = set() + for row in rows: + symbol = normalize_strategy_symbol(str(row.get("symbol") or "")) + side = str(row.get("side") or "").strip().lower() + if not symbol or side not in ("long", "short"): + continue + if symbol.upper() not in allowed: + continue + key = (symbol, side) + if key in seen: + continue + seen.add(key) + output.append(key) + return output + + +def _purge_flat_grace_sec() -> float: + """Guard window after an open/add fill during which a flat exchange snapshot must not purge L3.""" + try: + return max(0.0, float(os.getenv("POSITION_SYNC_PURGE_GRACE_SEC", "120"))) + except Exception: + return 120.0 + + +def _recent_open_fill_symbols(strategy_id: int, grace_sec: float) -> set[tuple[str, str]]: + """Return (symbol, side) legs that had an open/add fill within the grace window.""" + if grace_sec <= 0: + return set() + sid = int(strategy_id or 0) + if sid <= 0: + return set() + out: set[tuple[str, str]] = set() + try: + with get_db_connection() as db: + cur = db.cursor() + cur.execute( + """ + SELECT symbol, type + FROM qd_strategy_trades + WHERE strategy_id = %s + AND type IN ('open_long', 'open_short', 'add_long', 'add_short') + AND created_at >= NOW() - (%s * INTERVAL '1 second') + """, + (sid, grace_sec), + ) + rows = cur.fetchall() or [] + cur.close() + except Exception as exc: + logger.warning("[PositionSync] recent open fill lookup failed sid=%s: %s", sid, exc) + return out + for row in rows: + symbol = normalize_strategy_symbol(str(row.get("symbol") or "")) + if not symbol: + continue + t = str(row.get("type") or "").strip().lower() + side = "short" if "short" in t else "long" if "long" in t else "" + if side: + out.add((symbol.upper(), side)) + return out + + +def _purge_flat_strategy_positions_from_exchange( + *, + strategy_id: int, + strategy_config: Dict[str, Any], + exch_size: Dict[str, Dict[str, float]], + client: Any = None, + eps: float = 1e-12, +) -> int: + """ + Remove local L3 strategy legs only when a fresh exchange snapshot confirms + that the corresponding exchange leg is flat. + + Before deleting a flat local leg, write a ``close_*`` trade when the strategy + trade ledger still shows residual open size (native protection closes are + otherwise invisible to the UI). + + A short grace window protects against a transient flat snapshot racing a + just-executed open/add fill. + """ + try: + enabled = str(os.getenv("POSITION_SYNC_PURGE_FLAT_LEDGER", "true")).strip().lower() + except Exception: + enabled = "true" + if enabled in {"0", "false", "no", "off"}: + return 0 + + sid = int(strategy_id or 0) + if sid <= 0: + return 0 + allowed = strategy_allowed_symbols(strategy_config or {}) + if not allowed: + return 0 + + grace_sec = _purge_flat_grace_sec() + recent_open = _recent_open_fill_symbols(sid, grace_sec) + + deleted = 0 + for symbol, side in _local_strategy_position_legs(sid, allowed): + if lookup_exchange_side_qty(exch_size or {}, symbol, side) > eps: + continue + if (symbol.upper(), side) in recent_open: + logger.info( + "[PositionSync] Strategy %s keeps leg %s %s: open fill within grace window (%.0fs)", + sid, + symbol, + side, + grace_sec, + ) + continue + try: + from app.services.live_trading.external_flat_close import reconcile_external_flat_closes + + trade_id = reconcile_external_flat_closes( + strategy_id=sid, + symbol=symbol, + side=side, + client=client, + clear_local_position=True, + ) + if trade_id: + logger.warning( + "[PositionSync] Strategy %s recorded external flat close trade_id=%s for %s %s", + sid, + trade_id, + symbol, + side, + ) + except Exception as reconcile_err: + logger.warning( + "[PositionSync] Strategy %s failed external flat close reconcile for %s %s: %s", + sid, + symbol, + side, + reconcile_err, + ) + _delete_position(sid, symbol, side) + deleted += 1 + + # Heal cases where L3 was already purged but trade ledger still shows open qty. + try: + from app.services.live_trading.external_flat_close import ( + reconcile_external_flat_closes, + trade_side_net_qty, + ) + + for symbol in allowed: + sym = normalize_strategy_symbol(str(symbol or "")) + if not sym: + continue + for side in ("long", "short"): + if lookup_exchange_side_qty(exch_size or {}, sym, side) > eps: + continue + if trade_side_net_qty(sid, sym, side) <= eps: + continue + trade_id = reconcile_external_flat_closes( + strategy_id=sid, + symbol=sym, + side=side, + client=client, + clear_local_position=False, + ) + if trade_id: + logger.warning( + "[PositionSync] Strategy %s backfilled missing close trade_id=%s for %s %s", + sid, + trade_id, + sym, + side, + ) + except Exception as heal_err: + logger.warning( + "[PositionSync] Strategy %s missing-close heal failed: %s", + sid, + heal_err, + ) + return deleted + + + def _position_sync_fd_backoff_sec() -> float: try: return max(30.0, float(os.getenv("POSITION_SYNC_FD_BACKOFF_SEC", "90"))) @@ -229,6 +442,7 @@ def _sync_positions_best_effort(self, target_strategy_id: Optional[int] = None) exch_size: Dict[str, Dict[str, float]] = {} exch_entry_price: Dict[str, Dict[str, float]] = {} exch_inst_id: Dict[str, Dict[str, str]] = {} + client = None if cached_snap is not None: exch_size, exch_entry_price, exch_inst_id = cached_snap @@ -630,6 +844,37 @@ def _sync_positions_best_effort(self, target_strategy_id: Optional[int] = None) except Exception as l1_err: logger.warning("[PositionSync] L1 account sync failed key=%s: %s", cache_key, l1_err) + if client is None: + try: + client = create_client(exchange_config, market_type=market_type) + except Exception as client_err: + logger.debug( + "[PositionSync] Strategy %s could not create client for flat-close reconcile: %s", + sid, + client_err, + ) + client = None + + try: + purged = _purge_flat_strategy_positions_from_exchange( + strategy_id=int(sid), + strategy_config=sc, + exch_size=exch_size, + client=client, + ) + if purged: + logger.warning( + "[PositionSync] Strategy %s purged %s local leg(s) that are flat on exchange.", + sid, + purged, + ) + except Exception as purge_err: + logger.warning( + "[PositionSync] Strategy %s failed to purge flat local legs: %s", + sid, + purge_err, + ) + # [DEBUG] Log all normalized exchange keys for inspection logger.debug(f"[PositionSync] Strategy {sid} Exchange Keys: {list(exch_size.keys())}") @@ -646,10 +891,10 @@ def _sync_positions_best_effort(self, target_strategy_id: Optional[int] = None) else: logger.debug(f"[PositionSync] Strategy {sid} ({safe_cfg.get('exchange_id', 'unknown')}) has NO positions on exchange.") - # Keep exchange truth in L1 only. Strategy positions (L3) must be - # produced by that strategy's own fills, otherwise two live - # strategies sharing ETH/USDT would both inherit the same - # exchange account position. + # Keep non-flat exchange truth in L1 only. Strategy positions + # (L3) must be produced by that strategy's own fills; the only + # L3 mutation here is removing legs that the fresh exchange + # snapshot confirms are flat. except Exception as e: msg = str(e) if is_file_descriptor_exhausted(e): diff --git a/backend_api_python/app/services/pending_order_worker.py b/backend_api_python/app/services/pending_order_worker.py index 433f402ca..909d5ecd0 100644 --- a/backend_api_python/app/services/pending_order_worker.py +++ b/backend_api_python/app/services/pending_order_worker.py @@ -1340,6 +1340,12 @@ def _fetch_pending_orders(self, limit: int = 50) -> List[Dict[str, Any]]: UPDATE pending_orders SET status = 'pending', updated_at = NOW(), + attempts = CASE + WHEN attempts >= max_attempts + AND COALESCE(payload_json, '') LIKE '%%leverage_retry%%' + THEN GREATEST(0, max_attempts - 1) + ELSE attempts + END, dispatch_note = CASE WHEN dispatch_note IS NULL OR dispatch_note = '' THEN 'requeued_stale_processing' ELSE dispatch_note @@ -1347,7 +1353,10 @@ def _fetch_pending_orders(self, limit: int = 50) -> List[Dict[str, Any]]: WHERE status = 'processing' AND COALESCE(client_order_id, '') = '' AND (updated_at IS NULL OR updated_at < NOW() - INTERVAL '%s seconds') - AND (attempts < max_attempts) + AND ( + attempts < max_attempts + OR COALESCE(payload_json, '') LIKE '%%leverage_retry%%' + ) """, (stale_sec,), ) @@ -1373,7 +1382,21 @@ def _fetch_pending_orders(self, limit: int = 50) -> List[Dict[str, Any]]: ) rows = cur.fetchall() or [] cur.close() - return rows + now_ts = time.time() + eligible: List[Dict[str, Any]] = [] + for row in rows: + payload_json = row.get("payload_json") or "" + payload_obj: Dict[str, Any] = {} + if payload_json and isinstance(payload_json, str): + try: + payload_obj = json.loads(payload_json) or {} + except Exception: + payload_obj = {} + not_before = float(payload_obj.get("leverage_retry_not_before") or 0) + if not_before > now_ts: + continue + eligible.append(row) + return eligible except Exception as e: logger.warning(f"fetch_pending_orders failed: {e}") return [] @@ -1952,13 +1975,16 @@ def _execute_live_order(self, *, order_id: int, order_row: Dict[str, Any], paylo margin_mode=margin_mode, ) except Exception as e: - err = f"derivatives_account_configuration_failed:{e}" - logger.warning(f"live leverage set failed: pending_id={order_id}, strategy_id={strategy_id}, cfg={safe_cfg}, err={e}") - self._mark_failed(order_id=order_id, error=err) - _console_print(f"[worker] order rejected: strategy_id={strategy_id} pending_id={order_id} {err}") - _notify_live_best_effort(status="failed", error=err, amount_hint=amount, price_hint=ref_price) - append_strategy_log(strategy_id, "error", f"Leverage or margin-mode setup failed for {symbol}: {e}") - return + from app.services.pending_orders.leverage_retry import ( + handle_derivatives_configuration_error as _hl, + ) + if _hl( + error=e, order_id=order_id, strategy_id=strategy_id, symbol=str(symbol), + payload=payload, phases=phases, safe_cfg=safe_cfg, mark_failed=self._mark_failed, + append_log=append_strategy_log, console_print=_console_print, + notify=lambda **kw: _notify_live_best_effort(amount_hint=amount, price_hint=ref_price, **kw), + ): + return fills = FillAccumulator() diff --git a/backend_api_python/app/services/pending_orders/leverage_retry.py b/backend_api_python/app/services/pending_orders/leverage_retry.py new file mode 100644 index 000000000..28b161044 --- /dev/null +++ b/backend_api_python/app/services/pending_orders/leverage_retry.py @@ -0,0 +1,197 @@ +"""Recoverable OKX leverage / algo-order (59669) retry helpers. + +OKX may reject leverage or margin-mode changes while native algo / protection +orders exist on the instrument. Account configuration already cancels and +retries once; if that still fails, pending live orders should be requeued with +backoff instead of being marked permanently failed. +""" + +from __future__ import annotations + +import json +import time +from typing import Any, Callable, Dict, Optional + +from app.utils.db import get_db_connection +from app.utils.logger import get_logger + +logger = get_logger(__name__) + +MAX_LEVERAGE_RETRIES = 3 + + +def is_recoverable_leverage_error(error: Any) -> bool: + err_str = str(error or "") + lower = err_str.lower() + return ( + "59669" in err_str + or "Cancel cross-margin trailing" in err_str + or "algo orders" in lower + ) + + +def schedule_leverage_retry( + *, + order_id: int, + payload: Dict[str, Any], + attempt_count: int, + max_retries: int, + retry_delay_sec: int, + err_str: str, +) -> None: + """Requeue a live order blocked by recoverable OKX leverage errors.""" + merged_payload = dict(payload or {}) + merged_payload["leverage_retry_attempt"] = int(attempt_count) + merged_payload["leverage_retry_not_before"] = time.time() + float(retry_delay_sec) + note = f"leverage_retry:{attempt_count}/{max_retries} in {retry_delay_sec}s" + last_error = f"[auto-retry-{attempt_count}]{err_str}"[:2000] + try: + with get_db_connection() as db: + cur = db.cursor() + cur.execute( + """ + UPDATE pending_orders + SET status = 'pending', + attempts = GREATEST(0, COALESCE(attempts, 0) - 1), + last_error = %s, + dispatch_note = %s, + payload_json = %s, + updated_at = NOW() + WHERE id = %s + """, + ( + last_error, + note[:200], + json.dumps(merged_payload, ensure_ascii=False), + int(order_id), + ), + ) + db.commit() + cur.close() + except Exception as db_err: + logger.warning( + "Failed to schedule leverage retry pending_id=%s: %s", + order_id, + db_err, + ) + try: + with get_db_connection() as db: + cur = db.cursor() + cur.execute( + """ + UPDATE pending_orders + SET status = 'pending', + attempts = GREATEST(0, COALESCE(attempts, 0) - 1), + last_error = %s, + updated_at = NOW() + WHERE id = %s AND status = 'processing' + """, + (last_error, int(order_id)), + ) + db.commit() + cur.close() + except Exception: + logger.warning( + "Failed fallback requeue for leverage retry pending_id=%s", + order_id, + exc_info=True, + ) + + +def handle_derivatives_configuration_error( + *, + error: Exception, + order_id: int, + strategy_id: int, + symbol: str, + payload: Dict[str, Any], + phases: Dict[str, Any], + safe_cfg: Any, + mark_failed: Callable[..., None], + append_log: Callable[..., None], + notify: Optional[Callable[..., None]] = None, + console_print: Optional[Callable[[str], None]] = None, +) -> bool: + """Handle account-configuration failures for live orders. + + Returns True when the caller should stop processing the order (retry + scheduled or hard-failed). + """ + err_str = str(error) + if is_recoverable_leverage_error(error): + attempt_count = int(payload.get("leverage_retry_attempt") or 0) + 1 + if attempt_count <= MAX_LEVERAGE_RETRIES: + retry_delay_sec = attempt_count * 30 # 30s → 60s → 90s + phases["leverage_retry_attempt"] = attempt_count + phases["leverage_retry_pending"] = True + logger.warning( + "Leverage setup temporarily blocked by OKX algo orders " + "(attempt %s/%s, pending_id=%s, strategy_id=%s). " + "Scheduling retry in %ss...", + attempt_count, + MAX_LEVERAGE_RETRIES, + order_id, + strategy_id, + retry_delay_sec, + ) + append_log( + strategy_id, + "warning", + ( + f"Leverage setup temporarily blocked — " + f"retrying in {retry_delay_sec}s " + f"(attempt {attempt_count}/{MAX_LEVERAGE_RETRIES})" + ), + ) + schedule_leverage_retry( + order_id=order_id, + payload=payload, + attempt_count=attempt_count, + max_retries=MAX_LEVERAGE_RETRIES, + retry_delay_sec=retry_delay_sec, + err_str=err_str, + ) + return True + + err = f"derivatives_account_configuration_failed[max_retries_exceeded]:{error}" + logger.error( + "Leverage setup permanently failed after %s retries: " + "pending_id=%s, strategy_id=%s, err=%s", + MAX_LEVERAGE_RETRIES, + order_id, + strategy_id, + error, + ) + mark_failed(order_id=order_id, error=err) + append_log( + strategy_id, + "error", + ( + f"Leverage or margin-mode setup failed for {symbol} " + f"after {MAX_LEVERAGE_RETRIES} retries. " + "Restart the strategy to try again." + ), + ) + return True + + err = f"derivatives_account_configuration_failed:{error}" + logger.warning( + "live leverage set failed: pending_id=%s, strategy_id=%s, cfg=%s, err=%s", + order_id, + strategy_id, + safe_cfg, + error, + ) + mark_failed(order_id=order_id, error=err) + if console_print is not None: + console_print( + f"[worker] order rejected: strategy_id={strategy_id} pending_id={order_id} {err}" + ) + if notify is not None: + notify(status="failed", error=err) + append_log( + strategy_id, + "error", + f"Leverage or margin-mode setup failed for {symbol}: {error}", + ) + return True diff --git a/backend_api_python/app/utils/trade_close_reason.py b/backend_api_python/app/utils/trade_close_reason.py index b761f1892..d24d9f4b3 100644 --- a/backend_api_python/app/utils/trade_close_reason.py +++ b/backend_api_python/app/utils/trade_close_reason.py @@ -40,6 +40,12 @@ SERVER_TAKE_PROFIT = "server_take_profit" SERVER_TRAILING_STOP = "server_trailing_stop" +# Exchange-side closes discovered by position sync (native TP/SL, liq, etc.) +EXCHANGE_NATIVE_CLOSE = "exchange_native_close" +EXCHANGE_FLAT_RECONCILE = "exchange_flat_reconcile" +EXCHANGE_LIQUIDATION = "exchange_liquidation" +EXCHANGE_ADL = "exchange_adl" + # Indicator / generic script signal close (non-grid) INDICATOR_SIGNAL = "indicator_signal" @@ -71,6 +77,10 @@ SERVER_STOP_LOSS: "止损平仓", SERVER_TAKE_PROFIT: "止盈平仓", SERVER_TRAILING_STOP: "追踪止损平仓", + EXCHANGE_NATIVE_CLOSE: "交易所原生平仓", + EXCHANGE_FLAT_RECONCILE: "交易所空仓对账平仓", + EXCHANGE_LIQUIDATION: "交易所强平", + EXCHANGE_ADL: "交易所ADL减仓", INDICATOR_SIGNAL: "信号触发平仓", LEGACY_SIGNAL_TRIGGER: "信号触发平仓", } @@ -100,6 +110,10 @@ SERVER_STOP_LOSS: "Stop-loss close", SERVER_TAKE_PROFIT: "Take-profit close", SERVER_TRAILING_STOP: "Trailing stop close", + EXCHANGE_NATIVE_CLOSE: "Exchange native close", + EXCHANGE_FLAT_RECONCILE: "Exchange flat reconcile close", + EXCHANGE_LIQUIDATION: "Exchange liquidation", + EXCHANGE_ADL: "Exchange ADL close", INDICATOR_SIGNAL: "Signal close", LEGACY_SIGNAL_TRIGGER: "Signal close", } diff --git a/backend_api_python/tests/test_derivatives_account_configuration.py b/backend_api_python/tests/test_derivatives_account_configuration.py index d98eb1751..0ae2db780 100644 --- a/backend_api_python/tests/test_derivatives_account_configuration.py +++ b/backend_api_python/tests/test_derivatives_account_configuration.py @@ -55,6 +55,26 @@ def unchanged(*_args, **_kwargs): assert client.set_margin_mode("cross") is True +def test_reduce_only_swap_close_skips_account_configuration(): + from app.services.live_trading.account_configuration import ( + requires_derivatives_account_configuration, + ) + + assert requires_derivatives_account_configuration( + market_type="swap", reduce_only=True + ) is False + + +def test_swap_open_still_requires_account_configuration(): + from app.services.live_trading.account_configuration import ( + requires_derivatives_account_configuration, + ) + + assert requires_derivatives_account_configuration( + market_type="swap", reduce_only=False + ) is True + + def test_binance_rejects_leverage_above_api_limit_instead_of_clamping(): client = BinanceFuturesClient.__new__(BinanceFuturesClient) client._signed_request = lambda *_args, **_kwargs: pytest.fail("request must not be sent") diff --git a/backend_api_python/tests/test_external_flat_close.py b/backend_api_python/tests/test_external_flat_close.py new file mode 100644 index 000000000..97c5a8832 --- /dev/null +++ b/backend_api_python/tests/test_external_flat_close.py @@ -0,0 +1,169 @@ +"""Tests for exchange-native flat close trade backfill.""" + +from datetime import datetime, timezone + +from app.services.live_trading import external_flat_close as mod +from app.utils.trade_close_reason import EXCHANGE_NATIVE_CLOSE + + +class _FakeOkx: + def __init__(self, history=None, fills=None): + self.history = history or [] + self.fills = fills or [] + + def _signed_request(self, method, path, params=None): + if path.endswith("/account/positions-history"): + return {"data": self.history} + if path.endswith("/trade/fills-history"): + return {"data": self.fills} + return {"data": []} + + +def test_resolve_okx_external_close_matches_entry_and_fee(monkeypatch): + monkeypatch.setattr(mod, "OkxClient", _FakeOkx) + opened = datetime(2026, 9, 12, 8, 0, 0, tzinfo=timezone.utc) + closed_ms = int(datetime(2026, 9, 14, 2, 12, 23, tzinfo=timezone.utc).timestamp() * 1000) + client = _FakeOkx( + history=[ + { + "posSide": "short", + "direction": "long", + "openAvgPx": "77240", + "closeAvgPx": "77348.6", + "uTime": str(closed_ms), + "type": "2", + } + ], + fills=[ + { + "side": "buy", + "posSide": "short", + "fillPx": "77348.6", + "fee": "-0.00773486", + "feeCcy": "USDT", + "tradeId": "2900408350", + "ts": str(closed_ms), + } + ], + ) + # isinstance checks OkxClient; patch the class used in isinstance. + assert isinstance(client, _FakeOkx) + monkeypatch.setattr(mod, "OkxClient", _FakeOkx) + + snap = mod.resolve_okx_external_close( + client, + symbol="BTC/USDT", + side="short", + entry_price=77240.0, + opened_after=opened, + ) + assert snap is not None + assert abs(float(snap["price"]) - 77348.6) < 1e-9 + assert abs(float(snap["fee"]) - 0.00773486) < 1e-9 + assert snap["exchange_fill_id"] == "2900408350" + assert snap["close_reason"] == EXCHANGE_NATIVE_CLOSE + + +def test_record_external_flat_close_writes_trade_and_clears_local(monkeypatch): + recorded = {} + applied = {} + + monkeypatch.setattr(mod, "_already_recorded_fill", lambda *a, **k: False) + monkeypatch.setattr( + mod, + "_fetch_position", + lambda strategy_id, symbol, side: {"size": 0.0002, "entry_price": 77240.0}, + ) + + def _apply(**kwargs): + applied.update(kwargs) + return (-0.02172, None, 77240.0) + + monkeypatch.setattr(mod, "apply_fill_to_local_position", _apply) + + def _record(**kwargs): + recorded.update(kwargs) + return 71 + + monkeypatch.setattr(mod, "record_trade", _record) + monkeypatch.setattr(mod, "append_strategy_log", lambda *a, **k: None) + + trade_id = mod.record_external_flat_close( + strategy_id=32, + symbol="BTC/USDT", + side="short", + amount=0.0002, + entry_price=77240.0, + close_price=77348.6, + commission=0.00773486, + close_reason=EXCHANGE_NATIVE_CLOSE, + exchange_fill_id="2900408350", + clear_local_position=True, + ) + assert trade_id == 71 + assert applied["signal_type"] == "close_short" + assert recorded["trade_type"] == "close_short" + assert abs(float(recorded["profit"]) - ((77240.0 - 77348.6) * 0.0002)) < 1e-9 + assert recorded["close_reason"] == EXCHANGE_NATIVE_CLOSE + + +def test_reconcile_defers_without_exchange_snapshot(monkeypatch): + monkeypatch.setattr(mod, "trade_side_net_qty", lambda *a, **k: 0.0002) + monkeypatch.setattr( + mod, + "_last_open_trade", + lambda *a, **k: { + "price": 77240.0, + "created_at": datetime(2026, 9, 12, 8, 0, 0), + "credential_id": 5, + "inst_id": "BTC-USDT-SWAP", + "market_type": "swap", + }, + ) + monkeypatch.setattr(mod, "resolve_okx_external_close", lambda *a, **k: None) + monkeypatch.setattr( + mod, + "record_external_flat_close", + lambda **kwargs: (_ for _ in ()).throw(AssertionError("should defer")), + ) + + trade_id = mod.reconcile_external_flat_closes( + strategy_id=32, + symbol="BTC/USDT", + side="short", + client=None, + fallback_price=0.0, + clear_local_position=False, + ) + assert trade_id == 0 + + +def test_resolve_falls_back_to_fills_history(monkeypatch): + monkeypatch.setattr(mod, "OkxClient", _FakeOkx) + opened = datetime(2026, 9, 14, 16, 0, 0, tzinfo=timezone.utc) + closed_ms = int(datetime(2026, 9, 14, 22, 4, 0, tzinfo=timezone.utc).timestamp() * 1000) + client = _FakeOkx( + history=[], + fills=[ + { + "side": "sell", + "posSide": "long", + "fillPx": "78657.8", + "fee": "-0.00786578", + "feeCcy": "USDT", + "tradeId": "2919999999", + "ts": str(closed_ms), + } + ], + ) + snap = mod.resolve_okx_external_close( + client, + symbol="BTC/USDT", + side="long", + entry_price=78552.9, + opened_after=opened, + ) + assert snap is not None + assert abs(float(snap["price"]) - 78657.8) < 1e-9 + assert snap["source"] == "fills_history" + assert snap["close_reason"] == EXCHANGE_NATIVE_CLOSE diff --git a/backend_api_python/tests/test_okx_cancel_algo_payload.py b/backend_api_python/tests/test_okx_cancel_algo_payload.py new file mode 100644 index 000000000..750da6cf6 --- /dev/null +++ b/backend_api_python/tests/test_okx_cancel_algo_payload.py @@ -0,0 +1,62 @@ +from __future__ import annotations + +from app.services.live_trading.okx import OkxClient + + +def test_cancel_algo_order_sends_array_body_to_advance_endpoint_for_trailing(): + calls = [] + + def fake_signed_request(method, path, **kwargs): + calls.append({"method": method, "path": path, "json_body": kwargs.get("json_body")}) + return {"code": "0", "data": [{"algoId": "1", "sCode": "0"}]} + + client = OkxClient.__new__(OkxClient) + client._signed_request = fake_signed_request + + client.cancel_algo_order( + algo_id="3807971948987043840", + inst_id="BTC-USDT-SWAP", + ord_type="move_order_stop", + ) + + assert calls + assert calls[0]["path"] == "/api/v5/trade/cancel-advance-algos" + assert isinstance(calls[0]["json_body"], list) + assert calls[0]["json_body"] == [ + {"algoId": "3807971948987043840", "instId": "BTC-USDT-SWAP"} + ] + + +def test_cancel_all_algo_orders_batches_by_type(): + calls = [] + + client = OkxClient.__new__(OkxClient) + + def fake_get_algo_orders(*, ord_type: str = "", **_kwargs): + if ord_type != "move_order_stop": + return {"data": []} + return { + "data": [ + {"algoId": "1", "instId": "BTC-USDT-SWAP"}, + {"algoId": "2", "instId": "BTC-USDT-SWAP"}, + ] + } + + def fake_signed_request(method, path, **kwargs): + calls.append({"path": path, "json_body": kwargs.get("json_body")}) + body = kwargs.get("json_body") or [] + return { + "code": "0", + "data": [{"algoId": item["algoId"], "sCode": "0"} for item in body], + } + + client.get_algo_orders = fake_get_algo_orders + client._signed_request = fake_signed_request + + cancelled = client.cancel_all_algo_orders(inst_id="BTC-USDT-SWAP", inst_type="SWAP") + + assert cancelled == 2 + assert calls + assert calls[0]["path"] == "/api/v5/trade/cancel-advance-algos" + assert isinstance(calls[0]["json_body"], list) + assert len(calls[0]["json_body"]) == 2 diff --git a/backend_api_python/tests/test_okx_set_leverage_59669.py b/backend_api_python/tests/test_okx_set_leverage_59669.py new file mode 100644 index 000000000..09694a011 --- /dev/null +++ b/backend_api_python/tests/test_okx_set_leverage_59669.py @@ -0,0 +1,62 @@ +from __future__ import annotations + +from app.services.live_trading.base import LiveTradingError +from app.services.live_trading.okx import OkxClient + + +def _client(**overrides) -> OkxClient: + client = OkxClient.__new__(OkxClient) + client._lev_cache = {} + client._lev_cache_ttl_sec = 60.0 + for key, value in overrides.items(): + setattr(client, key, value) + return client + + +def test_okx_set_leverage_skips_when_position_already_at_target(): + client = _client( + _read_configured_leverage=lambda **kwargs: 1, + _signed_request=lambda *args, **kwargs: (_ for _ in ()).throw( + AssertionError("set-leverage should not be called") + ), + ) + + assert client.set_leverage(inst_id="BTC-USDT-SWAP", lever=1, mgn_mode="cross", pos_side="net") is True + + +def test_okx_set_leverage_skips_when_leverage_info_already_at_target(): + client = _client( + _read_effective_leverage=lambda **kwargs: None, + _read_configured_leverage=lambda **kwargs: 1, + _signed_request=lambda *args, **kwargs: (_ for _ in ()).throw( + AssertionError("set-leverage should not be called") + ), + ) + + assert client.set_leverage(inst_id="BTC-USDT-SWAP", lever=1, mgn_mode="cross", pos_side="net") is True + + +def test_okx_set_leverage_59669_accepts_existing_leverage(): + calls = {"set": 0} + + def fake_read_configured_leverage(**_kwargs): + # Keep attempting set-leverage until retries are exhausted, then + # report the target leverage as already configured on the venue. + if calls["set"] >= 4: + return 1 + return None + + def fake_signed_request(method, path, **kwargs): + if path == "/api/v5/account/set-leverage": + calls["set"] += 1 + raise LiveTradingError("OKX error: {'code': '59669', 'msg': 'blocked'}") + raise AssertionError(f"unexpected path: {path}") + + client = _client( + _read_configured_leverage=fake_read_configured_leverage, + _signed_request=fake_signed_request, + cancel_all_algo_orders=lambda **kwargs: 0, + ) + + assert client.set_leverage(inst_id="BTC-USDT-SWAP", lever=1, mgn_mode="cross", pos_side="net") is True + assert calls["set"] == 4 diff --git a/backend_api_python/tests/test_pending_order_leverage_retry.py b/backend_api_python/tests/test_pending_order_leverage_retry.py new file mode 100644 index 000000000..3846a2abb --- /dev/null +++ b/backend_api_python/tests/test_pending_order_leverage_retry.py @@ -0,0 +1,86 @@ +from __future__ import annotations + +import json +import time +from unittest.mock import MagicMock, patch + +from app.services.pending_order_worker import PendingOrderWorker +from app.services.pending_orders import leverage_retry as leverage_retry_mod + + +def test_schedule_leverage_retry_requeues_pending_order(): + payload = {"strategy_id": 11, "signal_type": "add_short"} + + with patch("app.services.pending_orders.leverage_retry.get_db_connection") as mock_conn: + db = MagicMock() + cur = MagicMock() + mock_conn.return_value.__enter__.return_value = db + db.cursor.return_value = cur + + leverage_retry_mod.schedule_leverage_retry( + order_id=14, + payload=payload, + attempt_count=2, + max_retries=3, + retry_delay_sec=60, + err_str="OKX 59669", + ) + + cur.execute.assert_called_once() + sql, params = cur.execute.call_args[0] + assert "UPDATE pending_orders" in sql + assert "status = 'pending'" in sql + assert params[0] == "[auto-retry-2]OKX 59669" + stored_payload = json.loads(params[2]) + assert stored_payload["leverage_retry_attempt"] == 2 + assert stored_payload["leverage_retry_not_before"] > time.time() + db.commit.assert_called_once() + + +def test_fetch_pending_orders_skips_leverage_retry_backoff(): + worker = PendingOrderWorker.__new__(PendingOrderWorker) + worker._stale_processing_sec = 0 + + future_payload = json.dumps({"leverage_retry_not_before": time.time() + 3600}) + ready_payload = json.dumps({"leverage_retry_not_before": time.time() - 1}) + + with patch("app.services.pending_order_worker.get_db_connection") as mock_conn: + db = MagicMock() + cur = MagicMock() + mock_conn.return_value.__enter__.return_value = db + db.cursor.return_value = cur + cur.fetchall.return_value = [ + {"id": 1, "payload_json": future_payload, "status": "pending"}, + {"id": 2, "payload_json": ready_payload, "status": "pending"}, + ] + + rows = worker._fetch_pending_orders(limit=10) + + assert len(rows) == 1 + assert rows[0]["id"] == 2 + + +def test_handle_derivatives_configuration_error_schedules_retry(): + phases: dict = {} + marked = [] + logs = [] + + with patch.object(leverage_retry_mod, "schedule_leverage_retry") as schedule: + handled = leverage_retry_mod.handle_derivatives_configuration_error( + error=RuntimeError("OKX 59669 Cancel algo orders"), + order_id=9, + strategy_id=3, + symbol="BTC/USDT", + payload={}, + phases=phases, + safe_cfg={"exchange_id": "okx"}, + mark_failed=lambda **kw: marked.append(kw), + append_log=lambda *args: logs.append(args), + ) + + assert handled is True + assert phases["leverage_retry_pending"] is True + assert phases["leverage_retry_attempt"] == 1 + schedule.assert_called_once() + assert marked == [] + assert logs and logs[0][1] == "warning" diff --git a/backend_api_python/tests/test_position_sync_flat_purge.py b/backend_api_python/tests/test_position_sync_flat_purge.py new file mode 100644 index 000000000..cb2fd39b9 --- /dev/null +++ b/backend_api_python/tests/test_position_sync_flat_purge.py @@ -0,0 +1,139 @@ +from app.services import pending_order_position_sync as sync_mod + + +def _no_recent_open_fills(monkeypatch): + monkeypatch.setattr( + sync_mod, + "_recent_open_fill_symbols", + lambda strategy_id, grace_sec: set(), + ) + + +def test_purge_flat_strategy_positions_deletes_exchange_flat_legs(monkeypatch): + deleted = [] + + _no_recent_open_fills(monkeypatch) + monkeypatch.setattr( + sync_mod, + "_delete_position", + lambda strategy_id, symbol, side: deleted.append((strategy_id, symbol, side)), + ) + monkeypatch.setattr( + sync_mod, + "_local_strategy_position_legs", + lambda strategy_id, allowed_symbols: [("BTC/USDT", "long"), ("BTC/USDT", "short")], + ) + + count = sync_mod._purge_flat_strategy_positions_from_exchange( + strategy_id=7, + strategy_config={"symbol": "BTC/USDT"}, + exch_size={}, + ) + + assert count == 2 + assert (7, "BTC/USDT", "long") in deleted + assert (7, "BTC/USDT", "short") in deleted + + +def test_purge_flat_strategy_positions_keeps_exchange_live_leg(monkeypatch): + deleted = [] + + _no_recent_open_fills(monkeypatch) + monkeypatch.setattr( + sync_mod, + "_delete_position", + lambda strategy_id, symbol, side: deleted.append((strategy_id, symbol, side)), + ) + monkeypatch.setattr( + sync_mod, + "_local_strategy_position_legs", + lambda strategy_id, allowed_symbols: [("BTC/USDT", "long"), ("BTC/USDT", "short")], + ) + + count = sync_mod._purge_flat_strategy_positions_from_exchange( + strategy_id=8, + strategy_config={"symbol": "BTC/USDT"}, + exch_size={"BTC/USDT": {"long": 0.01, "short": 0.0}}, + ) + + assert count == 1 + assert (8, "BTC/USDT", "long") not in deleted + assert (8, "BTC/USDT", "short") in deleted + + +def test_purge_flat_strategy_positions_can_be_disabled(monkeypatch): + deleted = [] + _no_recent_open_fills(monkeypatch) + monkeypatch.setenv("POSITION_SYNC_PURGE_FLAT_LEDGER", "false") + monkeypatch.setattr( + sync_mod, + "_delete_position", + lambda strategy_id, symbol, side: deleted.append((strategy_id, symbol, side)), + ) + monkeypatch.setattr( + sync_mod, + "_local_strategy_position_legs", + lambda strategy_id, allowed_symbols: [("ETH/USDT", "long")], + ) + + count = sync_mod._purge_flat_strategy_positions_from_exchange( + strategy_id=9, + strategy_config={"symbol": "ETH/USDT"}, + exch_size={}, + ) + + assert count == 0 + assert deleted == [] + + +def test_purge_flat_strategy_positions_noops_when_local_ledger_empty(monkeypatch): + deleted = [] + _no_recent_open_fills(monkeypatch) + monkeypatch.setattr( + sync_mod, + "_delete_position", + lambda strategy_id, symbol, side: deleted.append((strategy_id, symbol, side)), + ) + monkeypatch.setattr( + sync_mod, + "_local_strategy_position_legs", + lambda strategy_id, allowed_symbols: [], + ) + + count = sync_mod._purge_flat_strategy_positions_from_exchange( + strategy_id=10, + strategy_config={"symbol": "ETH/USDT"}, + exch_size={}, + ) + + assert count == 0 + assert deleted == [] + + +def test_purge_flat_strategy_positions_keeps_leg_within_grace_window(monkeypatch): + deleted = [] + monkeypatch.setenv("POSITION_SYNC_PURGE_GRACE_SEC", "120") + monkeypatch.setattr( + sync_mod, + "_delete_position", + lambda strategy_id, symbol, side: deleted.append((strategy_id, symbol, side)), + ) + monkeypatch.setattr( + sync_mod, + "_local_strategy_position_legs", + lambda strategy_id, allowed_symbols: [("BTC/USDT", "long")], + ) + monkeypatch.setattr( + sync_mod, + "_recent_open_fill_symbols", + lambda strategy_id, grace_sec: {("BTC/USDT", "long")}, + ) + + count = sync_mod._purge_flat_strategy_positions_from_exchange( + strategy_id=11, + strategy_config={"symbol": "BTC/USDT"}, + exch_size={}, + ) + + assert count == 0 + assert deleted == []