Files
Hermes-Skills/okx-auto-position/scripts/process_signal.py
T
mike 8b70ce3048 fix: 加仓/减仓标签显示正确
【bug】加仓/减仓信号被推为'自动开仓', 因为:
1. format_execution_result 写死 ' ... 自动开仓'
2. type_labels 字典没 'add'/'reduce' 键
3. parse_signal 没存 _raw, classify_signal 拿不到原始 text, fallback 'open'

【修复】
1. parse_signal 存 fields['_raw'] = text (classify 拿到原始 text)
2. type_labels 加 'add'/'reduce' 键
3. format_execution_result 根据 fields['signal_type'] 显示标签
4. main() 里 fields['signal_type'] = classify_signal(fields) 写回

【效果】下次加仓信号 → ' BTC 做多 🟩 4x  加仓' (不再误显示开仓)
减仓信号也加好了 (新代码)
2026-07-24 00:08:00 +08:00

914 lines
36 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/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, get_last_signal, check_signal_lock
# ─── 解析 ────────────────────────────────────────────────────────────────
def parse_signal(text):
"""从TG信号文本提取关键字段"""
fields = {'_raw': text} # 存原始文本, 供 classify_signal 判定
# 交易员 - 找"【交易员】"标签, 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:
# 实在找不到 trader 标签 — 默认 X聚合社区 (用户 TG 唯一信号源)
# 比 "unknown" 更有用, 用户能立刻知道源头
fields['trader'] = 'X聚合社区'
# 字段映射
extractors = {
'symbol': r'【币种】\s*[:]?\s*(\S+)',
'side': r'【方向】\s*[:]?\s*(做多|做空)',
'leverage':r'【杠杆】\s*[:]?\s*(\d+)',
'size': r'【仓位(?:大小)?】\s*[:]?\s*([\d,.]+)', # 兼容 【仓位】和【仓位大小】
'unit': r'【仓位(?:大小)?】\s*[:]?\s*[\d,.]+\s+(BTC|USDT|ETH|SOL|DOGE|张|USD)\b', # 严格匹配已知单位 (必需)
'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(',', '') if key != 'unit' else m.group(1)
# 清理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'
# 加仓信号
if '加仓' in text or '追仓' in text:
return 'add'
# 默认为新开仓
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', '?')
unit = fields.get('unit', '') # 仓位单位: BTC/USDT/张
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类加仓',
'add': ' 加仓',
'reduce': 'B类减仓',
'close': '平仓',
}
type_label = type_labels.get(signal_type, signal_type)
# 信号源仓位(只展示,不参与计算)
src_info = f"📊 {trader} {size} {unit} {symbol}(价值${value})← 信号源,非你的仓位"
# 仓位变化对比
current_size = 0
try:
current_size = float(fields.get('size', '0').replace(',', ''))
comparison = format_comparison(trader, symbol, current_size)
except:
comparison = ""
# 跟单建议: fetch 真实持仓 + 算比例
our_position = None
our_advice = ""
try:
import ccxt
from okx_position_advisor import load_credentials, create_exchange
creds = load_credentials()
ex = create_exchange(creds)
positions = ex.fetch_positions()
for p in positions:
if symbol in p.get('symbol', '') and p.get('contracts', 0) != 0:
our_position = p
break
if our_position and current_size > 0:
our_contracts = float(our_position.get('contracts', 0))
# 算大佬减仓比例
last_signal = get_last_signal(trader, symbol)
last_size = (last_signal or {}).get('trader_size', 0)
if last_size and last_size > 0:
trader_delta_pct = (current_size - last_size) / last_size * 100
if trader_delta_pct < -0.5: # 大佬减仓 > 0.5%
# 跟同比例
our_reduce = our_contracts * abs(trader_delta_pct) / 100
# 取整 (BTC min=0.01张, 其他min=1张)
if symbol == 'BTC':
our_reduce = max(0.01, round(our_reduce, 2))
else:
our_reduce = max(1, round(our_reduce))
our_advice = f"\n🎯 跟单建议: 大佬减 {abs(trader_delta_pct):.1f}%, 你跟减 {our_reduce} 张 ({our_contracts}{our_contracts - our_reduce:.2f})"
except Exception as e:
our_advice = f"\n⚠️ 跟单计算跳过: {e}"
# 交易员评分
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}
{our_advice}
{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"
# 1000PEPE 归一化为 PEPE (OKX 实际合约名)
if symbol_usdt == '1000PEPE':
inst_id = 'PEPE-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', {})
# 根据 signal_type 显示动作 (open=新开仓, add=加仓, close=平仓, reduce=减仓)
action_labels = {
'open': '新开仓',
'add': ' 加仓',
'reduce': '减仓',
'close': '平仓',
}
action_label = action_labels.get(fields.get('signal_type', 'open'), '自动开仓')
msg = f"""✅ {symbol} {side_cn} {emoji} {leverage}x {action_label}
📊 信号源: {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)
fields['signal_type'] = signal_type # 写回 fields, 给 format_execution_result 用
# 平仓信号走独立通道 — 不看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)
# 锁检查: 同币种只跟一个 trader (在 execute 之前检查)
side_en = 'long' if 'long' in str(side).lower() or side == '做多' else 'short'
allowed, lock_msg = check_signal_lock(symbol, trader, side_en)
if not allowed:
return f"🔒 {lock_msg}"
# 用户要求:所有信号自动执行,只有余额不足才跳过
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')
elif signal_type == 'add':
# 加仓信号: 走 execute_order 自动加 (advisor 内部会检查余额, 不足返回 0 张)
exec_result = execute_order(symbol, side, leverage, rec)
if exec_result and 'error' not in exec_result and rec.get('contracts', 0) > 0:
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_add')
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_add_insufficient')
elif signal_type == 'reduce':
# 减仓信号: 锁了同方向 → 按大佬减仓比例, 自动减我们的同向持仓
try:
import ccxt as _ccxt
from okx_position_advisor import load_credentials, create_exchange
creds = load_credentials()
ex = create_exchange(creds)
positions = ex.fetch_positions()
our_pos = next((p for p in positions if symbol in p.get('symbol', '') and p.get('contracts', 0) != 0), None)
if not our_pos:
msg = f"⏭️ 减仓信号: 你无 {symbol} 持仓, 跳过"
else:
our_contracts = float(our_pos.get('contracts', 0))
# 算大佬减仓比例
last = get_last_signal(trader, symbol)
trader_before = (last or {}).get('trader_size', 0) or 0
trader_after = float(fields.get('size', '0').replace(',', ''))
if trader_before <= 0:
msg = f"⚠️ 减仓信号: 大佬前仓位未知, 跳过 (你有 {our_contracts} 张)"
else:
delta_pct = (trader_before - trader_after) / trader_before
if delta_pct <= 0:
msg = f"⚠️ 减仓信号: 大佬实际是加仓 +{delta_pct*100:.1f}%, 跳过"
else:
# 按比例减
reduce_amt = our_contracts * delta_pct
if symbol == 'BTC':
reduce_amt = max(0.01, round(reduce_amt, 2))
else:
reduce_amt = max(1, round(reduce_amt))
# 如果减完 < 0.01 张, 全平
if symbol == 'BTC' and our_contracts - reduce_amt < 0.01:
reduce_amt = our_contracts # 全平
# 同向减仓: long → sell, short → buy
close_side = 'sell' if our_pos.get('side') == 'long' else 'buy'
order = ex.create_order(
symbol=f'{symbol}/USDT:USDT',
type='market',
side=close_side,
amount=reduce_amt,
params={'reduceOnly': True}
)
new_contracts = our_contracts - reduce_amt
msg = f"""✅ 减仓执行 | {symbol} {our_pos.get('side')} {delta_pct*100:.1f}%
📊 大佬: {trader_before:,.2f}{trader_after:,.2f}
💼 你的: {our_contracts}{new_contracts:.2f}
🔻 平 {reduce_amt} 张 (市价 {close_side})"""
except Exception as e:
msg = f"⚠️ 减仓失败: {e}"
_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_reduce')
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()