777 lines
28 KiB
Python
777 lines
28 KiB
Python
#!/usr/bin/env python3
|
||
"""
|
||
交易信号处理器(no_agent模式):
|
||
1. 解析TG信号文本
|
||
2. 调advisor脚本获取正确金额
|
||
3. 格式化含📐完整模板
|
||
4. 推QQ
|
||
5. 信号去重/合并
|
||
|
||
用法: python3 process_signal.py "信号文本"
|
||
或: echo "信号文本" | python3 process_signal.py
|
||
|
||
cron模式: 作为no_agent cron job的script使用
|
||
"""
|
||
import sys
|
||
import os
|
||
import re
|
||
import json
|
||
import subprocess
|
||
import sqlite3
|
||
import hashlib
|
||
from pathlib import Path
|
||
from datetime import datetime, timedelta
|
||
|
||
SKILL_DIR = Path.home() / ".hermes/skills/trading/okx-auto-position"
|
||
ADVISOR = SKILL_DIR / "scripts" / "okx_position_advisor.py"
|
||
QQ_PUSH = Path.home() / ".hermes/scripts/push_to_qq.sh"
|
||
SIGNAL_DB = Path.home() / ".hermes/trading/signal_history.db"
|
||
DEDUP_DB = Path.home() / ".hermes/trading/signal_dedup.db"
|
||
|
||
# Import signal tracker
|
||
sys.path.insert(0, str(SKILL_DIR / "scripts"))
|
||
from signal_tracker import format_comparison, record_signal as _tracker_record, record_confirmed, format_trader_rating
|
||
|
||
# ─── 解析 ────────────────────────────────────────────────────────────────
|
||
|
||
def parse_signal(text):
|
||
"""从TG信号文本提取关键字段"""
|
||
fields = {}
|
||
|
||
# 交易员 - 找"【交易员】"标签, 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 = {
|
||
'symbol': r'【币种】\s*[::]?\s*(\S+)',
|
||
'side': r'【方向】\s*[::]?\s*(做多|做空)',
|
||
'leverage':r'【杠杆】\s*[::]?\s*(\d+)',
|
||
'size': r'【仓位大小】\s*[::]?\s*([\d,.]+)',
|
||
'value': r'【仓位价值】\s*[::]?\s*\$?\s*([\d,.]+)',
|
||
'entry': r'【开仓价】\s*[::]?\s*([\d,.]+)',
|
||
'current': r'【当前价】\s*[::]?\s*([\d,.]+)',
|
||
'pnl': r'【未实现盈亏】\s*[::]?\s*([-\d,.]+)',
|
||
'margin': r'【保证金】\s*[::]?\s*\$?\s*([\d,.]+)',
|
||
}
|
||
|
||
for key, pattern in extractors.items():
|
||
m = re.search(pattern, text)
|
||
if m:
|
||
fields[key] = m.group(1).replace(',', '')
|
||
|
||
# 清理symbol
|
||
if 'symbol' in fields:
|
||
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
|
||
|
||
# 方向转英文
|
||
if fields.get('side', '').startswith('做多'):
|
||
fields['side_en'] = 'long'
|
||
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():
|
||
conn = sqlite3.connect(str(DEDUP_DB))
|
||
conn.execute("""
|
||
CREATE TABLE IF NOT EXISTS recent_signals (
|
||
id TEXT PRIMARY KEY,
|
||
symbol TEXT,
|
||
trader TEXT,
|
||
timestamp REAL,
|
||
raw_text TEXT
|
||
)
|
||
""")
|
||
conn.execute("""
|
||
CREATE TABLE IF NOT EXISTS processed (
|
||
msg_hash TEXT PRIMARY KEY,
|
||
processed_at REAL
|
||
)
|
||
""")
|
||
conn.commit()
|
||
return conn
|
||
|
||
def is_duplicate(conn, text, symbol, trader):
|
||
"""检查是否重复信号(同交易员同币种2分钟内)"""
|
||
msg_hash = hashlib.md5(text.encode()).hexdigest()
|
||
|
||
# 检查完全相同的消息
|
||
row = conn.execute(
|
||
"SELECT 1 FROM processed WHERE msg_hash = ?", (msg_hash,)
|
||
).fetchone()
|
||
if row:
|
||
return True
|
||
|
||
# 检查同交易员同币种2分钟内的信号
|
||
cutoff = datetime.now().timestamp() - 120 # 2分钟
|
||
row = conn.execute(
|
||
"""SELECT 1 FROM recent_signals
|
||
WHERE symbol = ? AND trader = ? AND timestamp > ?
|
||
ORDER BY timestamp DESC LIMIT 1""",
|
||
(symbol, trader, cutoff)
|
||
).fetchone()
|
||
|
||
return row is not None
|
||
|
||
def record_signal(conn, text, symbol, trader):
|
||
"""记录信号用于去重"""
|
||
msg_hash = hashlib.md5(text.encode()).hexdigest()
|
||
now = datetime.now().timestamp()
|
||
|
||
conn.execute(
|
||
"INSERT OR REPLACE INTO processed (msg_hash, processed_at) VALUES (?, ?)",
|
||
(msg_hash, now)
|
||
)
|
||
conn.execute(
|
||
"INSERT OR REPLACE INTO recent_signals (id, symbol, trader, timestamp, raw_text) VALUES (?, ?, ?, ?, ?)",
|
||
(msg_hash, symbol, trader, now, text[:500])
|
||
)
|
||
|
||
# 清理1小时前的记录
|
||
cutoff = now - 3600
|
||
conn.execute("DELETE FROM recent_signals WHERE timestamp < ?", (cutoff,))
|
||
conn.execute("DELETE FROM processed WHERE processed_at < ?", (cutoff,))
|
||
conn.commit()
|
||
|
||
# ─── Advisor ──────────────────────────────────────────────────────────────
|
||
|
||
def run_advisor(symbol, side, leverage):
|
||
"""调advisor脚本获取正确数据"""
|
||
cmd = [
|
||
'python3', str(ADVISOR),
|
||
'--symbol', symbol,
|
||
'--side', side,
|
||
'--leverage', str(leverage),
|
||
'--json'
|
||
]
|
||
|
||
try:
|
||
result = subprocess.run(
|
||
cmd, capture_output=True, text=True, timeout=30,
|
||
cwd=str(ADVISOR.parent)
|
||
)
|
||
if result.returncode == 0:
|
||
return json.loads(result.stdout)
|
||
else:
|
||
return {'error': result.stderr.strip()[:200]}
|
||
except subprocess.TimeoutExpired:
|
||
return {'error': 'advisor超时'}
|
||
except json.JSONDecodeError:
|
||
return {'error': 'advisor输出非JSON'}
|
||
except Exception as e:
|
||
return {'error': str(e)}
|
||
|
||
# ─── 分类 ────────────────────────────────────────────────────────────────
|
||
|
||
def classify_signal(fields):
|
||
"""判断信号类型:加仓/新开仓/减仓/平仓"""
|
||
text = fields.get('_raw', '')
|
||
|
||
# 平仓信号
|
||
if '平仓' in text or '止盈' in text or '止损' in text:
|
||
return 'close'
|
||
|
||
# 减仓信号
|
||
pnl = float(fields.get('pnl', '0').replace('+', ''))
|
||
if '减仓' in text or (pnl < 0 and '减' in text):
|
||
return 'reduce'
|
||
|
||
# 默认为新开仓或加仓(由advisor判断)
|
||
return 'open'
|
||
|
||
# ─── 格式化 ──────────────────────────────────────────────────────────────
|
||
|
||
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 '🟥'
|
||
leverage = fields.get('leverage', '10')
|
||
trader = fields.get('trader', '?')
|
||
size = fields.get('size', '?')
|
||
value = fields.get('value', '?')
|
||
entry_price = fields.get('entry', '?')
|
||
pnl_str = fields.get('pnl', '0')
|
||
pnl = float(pnl_str.replace('+', '')) if pnl_str else 0
|
||
current = rec.get('price', fields.get('current', '?'))
|
||
|
||
pnl_emoji = '🔥' if pnl > 0 else '🔴'
|
||
pnl_sign = '+' if pnl > 0 else ''
|
||
|
||
# 性价比
|
||
cc = rec.get('cost_check', {})
|
||
rr = cc.get('rr_ratio', rec.get('rr', 0))
|
||
profit = cc.get('profit_amount', rec.get('tp_pnl', 0))
|
||
fee = cc.get('fee_cost', 0)
|
||
fee_pct = cc.get('fee_pct', 0)
|
||
net = cc.get('net_profit', 0)
|
||
rating_emoji = cc.get('rating_emoji', '⚠️')
|
||
rating_text = cc.get('rating_text', '未知')
|
||
|
||
# 信号类型标签
|
||
type_labels = {
|
||
'open': '新开仓' if not fields.get('_is_add') else 'A类加仓',
|
||
'reduce': 'B类减仓',
|
||
'close': '平仓',
|
||
}
|
||
type_label = type_labels.get(signal_type, signal_type)
|
||
|
||
# 信号源仓位(只展示,不参与计算)
|
||
src_info = f"📊 {trader} {size} {symbol}(价值${value})← 信号源,非你的仓位"
|
||
|
||
# 仓位变化对比
|
||
try:
|
||
current_size = float(fields.get('size', '0').replace(',', ''))
|
||
comparison = format_comparison(trader, symbol, current_size)
|
||
except:
|
||
comparison = ""
|
||
|
||
# 交易员评分
|
||
try:
|
||
trader_rating = format_trader_rating(trader)
|
||
except:
|
||
trader_rating = ""
|
||
|
||
msg = f"""⚡ 跟单建议 | {symbol} {side_cn} {emoji} {leverage}x({type_label})
|
||
|
||
{src_info}
|
||
入场: ${_fmt(entry_price)} | 当前: ${_fmt(current)}
|
||
浮盈: {pnl_sign}{pnl:.0f} {pnl_emoji}
|
||
|
||
📊 仓位变化
|
||
{comparison}
|
||
|
||
{trader_rating}
|
||
|
||
📐 性价比
|
||
• 你的仓位: {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}
|
||
入场: ${_fmt(entry_price)} | 当前: ${_fmt(current)}
|
||
浮盈: {pnl_sign}{pnl:.0f} {pnl_emoji}
|
||
|
||
⚠️ 余额不足,无法开仓
|
||
• 可用: {_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):
|
||
"""推送到QQ"""
|
||
try:
|
||
result = subprocess.run(
|
||
['bash', str(QQ_PUSH), message],
|
||
capture_output=True, text=True, timeout=15
|
||
)
|
||
return result.returncode == 0
|
||
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):
|
||
"""执行开仓订单"""
|
||
cmd = [
|
||
'python3', str(ADVISOR),
|
||
'--symbol', symbol,
|
||
'--side', side,
|
||
'--leverage', str(leverage),
|
||
'--execute', '--json',
|
||
'--rec-json', json.dumps(rec)
|
||
]
|
||
|
||
try:
|
||
result = subprocess.run(
|
||
cmd, capture_output=True, text=True, timeout=30,
|
||
cwd=str(ADVISOR.parent)
|
||
)
|
||
if result.returncode == 0:
|
||
return json.loads(result.stdout)
|
||
else:
|
||
return {'error': result.stderr.strip()[:200]}
|
||
except Exception as e:
|
||
return {'error': str(e)}
|
||
|
||
def format_execution_result(fields, rec, exec_result):
|
||
"""格式化执行结果"""
|
||
symbol = fields.get('symbol', '?')
|
||
side_cn = fields.get('side', '做多')
|
||
emoji = '🟩' if fields.get('side_en') == 'long' else '🟥'
|
||
leverage = fields.get('leverage', '10')
|
||
trader = fields.get('trader', '?')
|
||
size = fields.get('size', '?')
|
||
value = fields.get('value', '?')
|
||
|
||
cc = rec.get('cost_check', {})
|
||
rr = cc.get('rr_ratio', rec.get('rr', 0))
|
||
profit = cc.get('profit_amount', rec.get('tp_pnl', 0))
|
||
fee = cc.get('fee_cost', 0)
|
||
fee_pct = cc.get('fee_pct', 0)
|
||
net = cc.get('net_profit', 0)
|
||
rating_emoji = cc.get('rating_emoji', '⚠️')
|
||
rating_text = cc.get('rating_text', '未知')
|
||
|
||
pos = exec_result.get('position', {})
|
||
algo = exec_result.get('algo', {})
|
||
|
||
msg = f"""✅ {symbol} {side_cn} {emoji} {leverage}x 自动开仓
|
||
|
||
📊 信号源: {trader} {size} {symbol}(价值${value})
|
||
|
||
📐 性价比检查
|
||
• 盈亏比: {rr}:1 ✅
|
||
• 盈利额: +{profit:.2f} USDT ✅
|
||
• 手续费: {fee:.2f} USDT ({fee_pct:.1f}%) ✅
|
||
• 净盈利: {net:.2f} USDT ✅
|
||
• 评级: {rating_emoji} {rating_text}
|
||
|
||
✅ 执行结果
|
||
• 入场: ${pos.get('entry', rec.get('price', '?'))}
|
||
• 仓位: {pos.get('contracts', rec.get('contracts', '?'))}张
|
||
• TP: ${algo.get('tp', rec.get('tp_price', '?'))}
|
||
• SL: ${algo.get('sl', rec.get('sl_price', '?'))}
|
||
• 强平: ${pos.get('liq', '?')}
|
||
|
||
━━━ 当前全部持仓 ━━━
|
||
(查询中...)"""
|
||
|
||
# 尝试获取当前全部持仓
|
||
try:
|
||
acct_cmd = ['python3', '-c', f'''
|
||
import sys
|
||
sys.path.insert(0, "{ADVISOR.parent}")
|
||
from okx_position_advisor import load_credentials, create_exchange, get_account_info
|
||
creds = load_credentials()
|
||
exchange = create_exchange(creds)
|
||
info = get_account_info(exchange)
|
||
print(f"Free: {{info['usdt_free']:.2f}}")
|
||
for p in info['positions']:
|
||
print(f" {{p['symbol']}}: {{p['contracts']}}张 UPL={{p['pnl']:.2f}}")
|
||
''']
|
||
acct_result = subprocess.run(acct_cmd, capture_output=True, text=True, timeout=15)
|
||
if acct_result.returncode == 0:
|
||
msg = msg.replace("(查询中...)", f"\n```\n{acct_result.stdout.strip()}\n```")
|
||
except:
|
||
pass
|
||
|
||
return msg
|
||
|
||
# ─── 主流程 ──────────────────────────────────────────────────────────────
|
||
|
||
def process_signal(text):
|
||
"""处理一条信号"""
|
||
# 解析
|
||
fields = parse_signal(text)
|
||
fields['_raw'] = text
|
||
|
||
if not fields.get('symbol') or not fields.get('side'):
|
||
return "⚠️ 无法解析信号"
|
||
|
||
symbol = fields['symbol']
|
||
side = fields['side_en']
|
||
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 "⏭️ 重复信号,跳过"
|
||
|
||
# 平仓信号 → 已在上方独立处理,这里不再重复
|
||
|
||
# 调advisor
|
||
rec = run_advisor(symbol, side, leverage)
|
||
|
||
if 'error' in rec:
|
||
record_signal(dedup_conn, text, symbol, trader)
|
||
dedup_conn.close()
|
||
return f"⚠️ advisor错误: {rec['error']}"
|
||
|
||
# 性价比检查
|
||
cc = rec.get('cost_check', {})
|
||
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 = rec.get('contracts', 0) > 0 # 有可开张数=自动执行
|
||
|
||
if auto_execute and signal_type == 'open':
|
||
# 性价比高 + 新开仓 → 自动执行
|
||
exec_result = execute_order(symbol, side, leverage, rec)
|
||
if exec_result and 'error' not in exec_result:
|
||
msg = format_execution_result(fields, rec, exec_result)
|
||
_tracker_record(trader=trader, symbol=symbol, side=side,
|
||
leverage=int(leverage) if leverage else 10,
|
||
trader_size=float(fields.get('size', '0').replace(',', '')),
|
||
trader_entry=float(fields.get('entry', '0').replace(',', '')),
|
||
trader_pnl=float(fields.get('pnl', '0').replace(',', '')),
|
||
raw_text=text, outcome='auto_executed')
|
||
else:
|
||
# 执行失败,降级为确认模式
|
||
auto_execute = False
|
||
msg = format_message(fields, rec, signal_type)
|
||
_tracker_record(trader=trader, symbol=symbol, side=side,
|
||
leverage=int(leverage) if leverage else 10,
|
||
trader_size=float(fields.get('size', '0').replace(',', '')),
|
||
trader_entry=float(fields.get('entry', '0').replace(',', '')),
|
||
trader_pnl=float(fields.get('pnl', '0').replace(',', '')),
|
||
raw_text=text, outcome='pushed')
|
||
else:
|
||
# 需要确认或减仓信号
|
||
msg = format_message(fields, rec, signal_type)
|
||
_tracker_record(trader=trader, symbol=symbol, side=side,
|
||
leverage=int(leverage) if leverage else 10,
|
||
trader_size=float(fields.get('size', '0').replace(',', '')),
|
||
trader_entry=float(fields.get('entry', '0').replace(',', '')),
|
||
trader_pnl=float(fields.get('pnl', '0').replace(',', '')),
|
||
raw_text=text, outcome='pushed')
|
||
|
||
# 记录去重
|
||
record_signal(dedup_conn, text, symbol, trader)
|
||
dedup_conn.close()
|
||
|
||
# 推送
|
||
success = push_to_qq(msg)
|
||
if success:
|
||
return f"✅ 已推送 | {symbol} {side} {leverage}x | {rec['contracts']}张 | 性价比{rec.get('cost_check', {}).get('rating_text', '?').replace('性价比', '')}"
|
||
else:
|
||
return f"❌ 推送失败"
|
||
|
||
def main():
|
||
if len(sys.argv) > 1:
|
||
text = ' '.join(sys.argv[1:])
|
||
else:
|
||
text = sys.stdin.read()
|
||
|
||
if not text.strip():
|
||
print("用法: python3 process_signal.py '信号文本'")
|
||
print("或: echo '信号文本' | python3 process_signal.py")
|
||
return
|
||
|
||
result = process_signal(text)
|
||
print(result)
|
||
|
||
if __name__ == '__main__':
|
||
main()
|