diff --git a/okx-auto-position/scripts/process_signal.py b/okx-auto-position/scripts/process_signal.py index 81addf0..f6bbf74 100644 --- a/okx-auto-position/scripts/process_signal.py +++ b/okx-auto-position/scripts/process_signal.py @@ -13,6 +13,7 @@ cron模式: 作为no_agent cron job的script使用 """ import sys +import os import re import json import subprocess @@ -37,10 +38,25 @@ def parse_signal(text): """从TG信号文本提取关键字段""" fields = {} - # 交易员 - m = re.search(r'【([^】]{1,20})】', text) - if m: - fields['trader'] = m.group(1) + # 交易员 - 找"【交易员】"标签, fallback "👉 跟单就选 X",再 fallback 第一个非字段名的方括号 + m_trader = re.search(r'【交易员】\s*[::]?\s*([^【\n]{1,20})', text) + if m_trader: + fields['trader'] = m_trader.group(1).strip() + else: + # Fallback: 👉 跟单就选 X (这是真 trader 来源) + m_follow = re.search(r'👉\s*跟单就选\s*(\S+)', text) + if m_follow: + fields['trader'] = m_follow.group(1).strip() + else: + # 最后 fallback: 第一个【xx】但跳过字段名 + FIELD_NAMES = {'币种', '方向', '杠杆', '仓位大小', '仓位价值', '开仓价', + '当前价', '未实现盈亏', '收益额', '持仓量', '强平价', '数量'} + m_first = re.search(r'【([^】]{1,20})】', text) + if m_first and m_first.group(1) not in FIELD_NAMES: + fields['trader'] = m_first.group(1).strip() + else: + # 实在找不到 → 标 "unknown",不强行填空 + fields['trader'] = 'unknown' # 字段映射 extractors = { @@ -62,8 +78,12 @@ def parse_signal(text): # 清理symbol if 'symbol' in fields: - sym = fields['symbol'] - sym = re.sub(r'\|.*$', '', sym) # 去掉 |永续|10x + sym_raw = fields['symbol'] + # 提取 |Nx 杠杆 (e.g. SKHYUSDT|永续|5x → 5) + m_lev = re.search(r'\|(\d+)\s*x?$', sym_raw) + if m_lev and 'leverage' not in fields: + fields['leverage'] = m_lev.group(1) + sym = re.sub(r'\|.*$', '', sym_raw) # 去掉 |永续|10x sym = sym.replace('USDT', '').strip() fields['symbol'] = sym @@ -73,8 +93,35 @@ def parse_signal(text): else: fields['side_en'] = 'short' + # 信号首次发出时间 (forwarder 加的 ⏱信号时间: 2026-07-08 18:00:00) + m = re.search(r'⏱信号时间[::]\s*(\d{4}-\d{2}-\d{2}\s+\d{2}:\d{2}:\d{2})', text) + if m: + try: + fields['signal_time'] = datetime.strptime(m.group(1), '%Y-%m-%d %H:%M:%S') + except ValueError: + pass + return fields +SIGNAL_FRESH_MINUTES = 30 # 30 分钟内算新鲜;>=30 算过期(少误跟) + +def is_signal_stale(fields): + """信号是否过期(>= 30 分钟)。 + 没有时间戳的按"新鲜"处理(不阻断旧信号源)。 + """ + st = fields.get('signal_time') + if not st: + return False + age = datetime.now() - st + return age.total_seconds() >= SIGNAL_FRESH_MINUTES * 60 + +def format_age_minutes(fields): + """信号已发出多久(分钟)。""" + st = fields.get('signal_time') + if not st: + return "?" + return int((datetime.now() - st).total_seconds() // 60) + # ─── 去重 ──────────────────────────────────────────────────────────────── def init_dedup_db(): @@ -192,6 +239,13 @@ def format_message(fields, rec, signal_type): if 'error' in rec: return f"⚠️ advisor错误: {rec['error']}" + def _fmt(x, n=4): + """格式化数字: 字符串保留原样, 数字 round 到 n 位.""" + try: + return f"{float(x):.{n}f}" + except (ValueError, TypeError): + return str(x) + symbol = fields.get('symbol', '?') side_cn = fields.get('side', '做多') emoji = '🟩' if fields.get('side_en') == 'long' else '🟥' @@ -244,7 +298,7 @@ def format_message(fields, rec, signal_type): msg = f"""⚡ 跟单建议 | {symbol} {side_cn} {emoji} {leverage}x({type_label}) {src_info} -入场: ${entry_price} | 当前: ${current} +入场: ${_fmt(entry_price)} | 当前: ${_fmt(current)} 浮盈: {pnl_sign}{pnl:.0f} {pnl_emoji} 📊 仓位变化 @@ -252,34 +306,52 @@ def format_message(fields, rec, signal_type): {trader_rating} -📐 性价比检查(基于你的推荐仓位) -• 你的仓位: {rec['contracts']}张(保证金{rec['margin']:.2f} USDT) -• 盈亏比: {rr}:1 {'✅' if rr >= 2 else '⚠️' if rr >= 1.5 else '❌'} -• 盈利额: +{profit:.2f} USDT {'✅' if profit >= 10 else '❌ <10U保底'} -• 手续费: {fee:.2f} USDT ({fee_pct:.1f}%) {'✅' if fee_pct < 5 else '❌'} -• 净盈利: {net:.2f} USDT {'✅' if net >= 10 else '❌'} -• 评级: {rating_emoji} {rating_text} -• SL: ${rec['sl_price']}(-{rec['sl_pct']:.1f}%) -• TP: ${rec['tp_price']}(+{rec['tp_pct']:.1f}%) - -回复 Y 确认跟单 / N 取消""" +📐 性价比 +• 你的仓位: {rec['contracts']}张(保证金{_fmt(rec['margin'])} USDT) +• SL: ${_fmt(rec['sl_price'])} → 预亏 -{_fmt(rec.get('sl_pnl', 0))} USDT (保证金-{rec.get('sl_pnl', 0)/max(rec.get('margin', 1), 0.01)*100:.0f}%) +• TP: ${_fmt(rec['tp_price'])} → 预盈 +{_fmt(rec.get('tp_pnl', 0))} USDT (保证金+{rec.get('tp_pnl', 0)/max(rec.get('margin', 1), 0.01)*100:.0f}%) +• 盈亏比: {rr}:1 {rating_emoji} {rating_text}""" # 如果余额不足,替换跟单方案 if rec.get('contracts', 0) == 0: msg = f"""⚡ 跟单建议 | {symbol} {side_cn} {emoji} {leverage}x({type_label}) {src_info} -入场: ${entry_price} | 当前: ${current} +入场: ${_fmt(entry_price)} | 当前: ${_fmt(current)} 浮盈: {pnl_sign}{pnl:.0f} {pnl_emoji} ⚠️ 余额不足,无法开仓 -• 可用: {rec.get('acct_free', 0):.2f} USDT -• 需要: ~{rec.get('margin', 0):.2f} USDT +• 可用: {_fmt(rec.get('acct_free', 0))} USDT +• 需要: ~{_fmt(rec.get('margin', 0))} USDT 💡 建议:等待其他仓位止盈释放保证金""" return msg +def format_stale_message(fields, age_minutes): + """过期信号提醒(不发 advisor 分析结果,只提示)。""" + symbol = fields.get('symbol', '?') + side_cn = fields.get('side', '?') + emoji = '🟩' if fields.get('side_en') == 'long' else '🟥' + leverage = fields.get('leverage', '?') + trader = fields.get('trader', '?') + entry_price = fields.get('entry', '?') + st = fields.get('signal_time') + + return f"""⏰ 信号已过期 | {symbol} {side_cn} {emoji} {leverage}x + +{trader} 信号源(仅展示,非你的仓位) + +信号首次发出: {st.strftime('%Y-%m-%d %H:%M:%S') if st else '?'} +已过去: {age_minutes} 分钟 (阈值 {SIGNAL_FRESH_MINUTES} 分钟) + +入场: ${entry_price} + +⚠️ 信号过期,谨慎跟单 +• 行情可能已经反转 +• 价格/仓位快照与当前不一致 +• 如需跟单请用实时数据重新评估""" + # ─── 推送 ──────────────────────────────────────────────────────────────── def push_to_qq(message): @@ -293,6 +365,196 @@ def push_to_qq(message): except: return False +# ─── 平仓 (raw REST, 绕 ccxt load_markets) ──────────────────────────────── + +def _okx_raw_request(method, path, params=None, body=None, timeout=15): + """OKX raw REST 调用 (避 ccxt fetch_balance→load_markets 超时)。 + + 签名规则 (按 ccxt/okx.py sign()): + auth = timestamp + method + request_path + if GET and query: auth += '?' + urlencode(query) + else: auth += json.dumps(body) + """ + import hmac, hashlib, base64, urllib.parse + creds = {} + with open(os.path.expanduser('~/.bashrc')) as f: + for line in f: + line = line.strip() + if line.startswith('export OKX_'): + k, v = line.replace('export ', '').split('=', 1) + creds[k] = v.strip().strip('"').strip("'") + for k, v in creds.items(): + if '${' not in v: + os.environ[k] = v + import re + for k, v in creds.items(): + if '${' in v: + os.environ[k] = re.sub(r'\$\{(\w+)\}', lambda m: os.environ.get(m.group(1), ''), v) + + ts = datetime.utcnow().strftime('%Y-%m-%dT%H:%M:%S.') + f"{datetime.utcnow().microsecond // 1000:03d}Z" + body_str = json.dumps(body) if body else '' + + # 构建签名 message + auth = ts + method.upper() + path + if method.upper() == 'GET': + if params: + # OKX: query 字符串按字典序排序后用 ? 拼接到 auth + sorted_q = '&'.join(f"{k}={urllib.parse.quote_plus(str(v), safe='')}" for k, v in sorted(params.items())) + auth += '?' + sorted_q + query = '?' + sorted_q + else: + query = '' + else: # POST + auth += body_str + query = '' + + sig = base64.b64encode(hmac.new(os.environ['OKX_SECRET'].encode(), auth.encode(), hashlib.sha256).digest()).decode() + headers = { + 'OK-ACCESS-KEY': os.environ['OKX_API_KEY'], + 'OK-ACCESS-SIGN': sig, + 'OK-ACCESS-TIMESTAMP': ts, + 'OK-ACCESS-PASSPHRASE': os.environ['OKX_PASSPHRASE'], + 'Content-Type': 'application/json', + } + + import requests + proxies = {'http': 'http://127.0.0.1:7890', 'https': 'http://127.0.0.1:7890'} + url = f'https://www.okx.com{path}{query}' + if method.upper() == 'GET': + r = requests.get(url, headers=headers, proxies=proxies, timeout=timeout) + else: + r = requests.post(url, data=body_str, headers=headers, proxies=proxies, timeout=timeout) + try: + return r.json() + except Exception: + return {'code': '-1', 'msg': f'非JSON响应: {r.text[:200]}'} + +def close_position_raw(symbol_usdt, signal_side): + """用 raw REST 平掉同币种同方向持仓。 + + Args: + symbol_usdt: 'ETH' / 'SKHY' / 'BTC' (base currency) + signal_side: 'long' / 'short' + + Returns: + dict: {action, pos_before, pos_after, pnl, message} + """ + inst_id = f"{symbol_usdt}-USDT-SWAP" + + # 1. 查当前持仓 + resp = _okx_raw_request('GET', '/api/v5/account/positions', {'instId': inst_id}) + if resp.get('code') != '0': + return {'action': 'error', 'message': f"查持仓失败: {resp.get('msg')}"} + + pos_before = None + for p in resp.get('data', []): + pos_size = float(p.get('pos', '0') or 0) + if pos_size > 0: + pos_before = { + 'pos': pos_size, + 'side': 'long', + 'avgPx': float(p.get('avgPx', '0') or 0), + 'upl': float(p.get('upl', '0') or 0), + 'lever': p.get('lever', '?'), + } + break + elif pos_size < 0: + pos_before = { + 'pos': abs(pos_size), + 'side': 'short', + 'avgPx': float(p.get('avgPx', '0') or 0), + 'upl': float(p.get('upl', '0') or 0), + 'lever': p.get('lever', '?'), + } + break + + if not pos_before: + return {'action': 'none', 'message': f'无 {symbol_usdt} 持仓'} + + # 2. 方向二次校验 + if pos_before['side'] != signal_side: + return { + 'action': 'skip', + 'pos_before': pos_before, + 'message': f"方向错位: 信号说{signal_side}但你持仓是{pos_before['side']},不动" + } + + # 3. 取最新价做参考 + tk = _okx_raw_request('GET', '/api/v5/market/ticker', {'instId': inst_id}) + current_px = None + if isinstance(tk, dict) and tk.get('code') == '0' and tk.get('data'): + try: + data0 = tk['data'][0] + # 用 dict.get 链避免 LSP 类型推断 (raw JSON 实际是 dict) + if isinstance(data0, dict): + last_val = data0.get('last') + if last_val is not None: + current_px = float(last_val) + except (KeyError, ValueError, TypeError): + pass + + # 4. 市价全平 reduceOnly + close_side = 'sell' if signal_side == 'long' else 'buy' + body = { + 'instId': inst_id, + 'tdMode': 'cross', + 'side': close_side, + 'posSide': 'net', + 'ordType': 'market', + 'sz': str(pos_before['pos']), + 'reduceOnly': True, + } + order_resp = _okx_raw_request('POST', '/api/v5/trade/order', body=body) + + if order_resp.get('code') == '0': + return { + 'action': 'closed', + 'pos_before': pos_before, + 'current_px': current_px, + 'order_id': order_resp.get('data', [{}])[0].get('ordId'), + 'message': '已市价全平', + } + else: + return { + 'action': 'error', + 'pos_before': pos_before, + 'message': f"下单失败: {order_resp.get('msg', order_resp)}" + } + +def format_close_message(fields, close_result): + """平仓信号处理结果推送 (精简版:v4.5.2 禁止过度分析)。""" + symbol = fields.get('symbol', '?') + side_cn = fields.get('side', '?') + emoji = '🟩' if fields.get('side_en') == 'long' else '🟥' + leverage = fields.get('leverage', '?') + trader = fields.get('trader', '?') + + action = close_result['action'] + + if action == 'closed': + pos = close_result['pos_before'] + u = pos['upl'] + u_emoji = '🔥' if u >= 0 else '🔴' + u_sign = '+' if u >= 0 else '' + return f"""✅ {symbol} {side_cn} {emoji} {leverage}x 已市价全平 | {trader} + +{pos['side']} {pos['pos']}张 | 浮盈 {u_sign}{u:.2f} USDT {u_emoji}""" + elif action == 'skip': + pos = close_result['pos_before'] + return f"""⏭️ {symbol} 平仓信号跳过 | {trader} + +你有反向持仓 {pos['side']} {pos['pos']}张 @ {pos['avgPx']:.2f} (浮盈 {pos['upl']:+.2f})""" + elif action == 'none': + return f"""ℹ️ {symbol} {side_cn} 平仓信号 | {trader} + +无持仓可平""" + else: # error + return f"""⚠️ {symbol} {side_cn} 平仓信号处理失败 | {trader} + +{close_result.get('message', '未知错误')} + +需手动处理""" + # ─── 执行订单 ──────────────────────────────────────────────────────────── def execute_order(symbol, side, leverage, rec): @@ -398,26 +660,45 @@ def process_signal(text): leverage = fields.get('leverage', '10') trader = fields.get('trader', '未知') - # 去重 + # 过期检查(30 分钟阈值,由 forwarder 注入的 ⏱信号时间 决定) + if is_signal_stale(fields): + dedup_conn = init_dedup_db() + record_signal(dedup_conn, text, symbol, trader) + dedup_conn.close() + age = format_age_minutes(fields) + msg = format_stale_message(fields, age) + push_to_qq(msg) + return f"⏰ 信号已过期 ({age}min) | 已推过期提醒" + + # 分类先于去重(让 close 信号绕过2分钟去重,因为平仓是必须执行的) + signal_type = classify_signal(fields) + + # 平仓信号走独立通道 — 不看2分钟窗口,只看 raw_text hash 是否完全重复 + # 修 2026-07-08 bug: 同币种同交易员的"减仓→平仓"紧跟信号被 dedup 误跳, + # 导致平仓规则从未触发,持仓长期不平 + if signal_type == 'close': + dedup_conn = init_dedup_db() + # close 信号只看 raw_text 是否完全相同(text 内含收益额+标记价,天然唯一) + msg_hash = hashlib.md5(text.encode()).hexdigest() + if dedup_conn.execute("SELECT 1 FROM processed WHERE msg_hash = ?", (msg_hash,)).fetchone(): + dedup_conn.close() + return "⏭️ 平仓信号完全重复跳过" + symbol_usdt = fields['symbol'] + signal_side_en = fields.get('side_en', 'long') + close_result = close_position_raw(symbol_usdt, signal_side_en) + msg = format_close_message(fields, close_result) + record_signal(dedup_conn, text, symbol, trader) + dedup_conn.close() + push_to_qq(msg) + return f"✅ 平仓处理: {close_result['action']} | {symbol_usdt} {signal_side_en}" + + # 其他信号(open/reduce): 才走2分钟窗口去重 dedup_conn = init_dedup_db() if is_duplicate(dedup_conn, text, symbol, trader): dedup_conn.close() return "⏭️ 重复信号,跳过" - # 分类 - signal_type = classify_signal(fields) - - # 平仓信号直接推送 - if signal_type == 'close': - msg = f"""🔔 {trader} {symbol}平仓提醒 -{text[text.find("入场"):text.find("回复")].strip() if "入场" in text else "详情见原始信号"} - -💡 操作建议 -• 若已跟单{symbol},建议同步止盈/止损""" - record_signal(dedup_conn, text, symbol, trader) - dedup_conn.close() - push_to_qq(msg) - return "✅ 平仓信号已推送" + # 平仓信号 → 已在上方独立处理,这里不再重复 # 调advisor rec = run_advisor(symbol, side, leverage) @@ -432,7 +713,8 @@ def process_signal(text): rr = cc.get('rr_ratio', rec.get('rr', 0)) profit = cc.get('profit_amount', rec.get('tp_pnl', 0)) fee_pct = cc.get('fee_pct', 0) - auto_execute = cc.get('auto_execute', False) or (rr >= 2 and fee_pct < 5 and profit >= 10) + # 用户要求:所有信号自动执行,只有余额不足才跳过 + auto_execute = rec.get('contracts', 0) > 0 # 有可开张数=自动执行 if auto_execute and signal_type == 'open': # 性价比高 + 新开仓 → 自动执行 @@ -472,7 +754,7 @@ def process_signal(text): # 推送 success = push_to_qq(msg) if success: - return f"✅ 已推送 | {symbol} {side} {leverage}x | {rec['contracts']}张 | 性价比{rec.get('cost_check', {}).get('rating_text', '?')}" + return f"✅ 已推送 | {symbol} {side} {leverage}x | {rec['contracts']}张 | 性价比{rec.get('cost_check', {}).get('rating_text', '?').replace('性价比', '')}" else: return f"❌ 推送失败"