Files
Hermes-Skills/okx-auto-position/scripts/signal_queue.py
T

123 lines
4.1 KiB
Python

#!/usr/bin/env python3
"""信号队列 - 先入库再处理,防止信号丢失"""
import json, os, time, sys, sqlite3
from datetime import datetime, timezone, timedelta
DB_PATH = os.path.expanduser("~/.hermes/trading/signal_queue.db")
def get_db():
os.makedirs(os.path.dirname(DB_PATH), exist_ok=True)
conn = sqlite3.connect(DB_PATH)
conn.execute("""
CREATE TABLE IF NOT EXISTS queue (
id INTEGER PRIMARY KEY AUTOINCREMENT,
raw_text TEXT NOT NULL,
status TEXT DEFAULT 'pending', -- pending/processing/done/failed
created_at TEXT DEFAULT (datetime('now')),
processed_at TEXT,
result TEXT,
error TEXT,
retries INTEGER DEFAULT 0
)
""")
conn.commit()
return conn
def enqueue(raw_text):
"""信号入队"""
conn = get_db()
conn.execute("INSERT INTO queue (raw_text, status) VALUES (?, 'pending')", (raw_text,))
conn.commit()
row_id = conn.execute("SELECT last_insert_rowid()").fetchone()[0]
conn.close()
return row_id
def get_pending(limit=10):
"""获取待处理信号"""
conn = get_db()
rows = conn.execute(
"SELECT id, raw_text, retries FROM queue WHERE status IN ('pending','failed') AND retries < 3 ORDER BY id LIMIT ?",
(limit,)
).fetchall()
conn.close()
return rows
def mark_processing(row_id):
conn = get_db()
conn.execute("UPDATE queue SET status='processing' WHERE id=?", (row_id,))
conn.commit()
conn.close()
def mark_done(row_id, result=""):
conn = get_db()
conn.execute("UPDATE queue SET status='done', processed_at=datetime('now'), result=? WHERE id=?",
(result[:500], row_id))
conn.commit()
conn.close()
def mark_failed(row_id, error=""):
conn = get_db()
conn.execute("UPDATE queue SET status='failed', processed_at=datetime('now'), error=?, retries=retries+1 WHERE id=?",
(error[:500], row_id))
conn.commit()
conn.close()
def get_stats():
conn = get_db()
stats = {}
for status in ['pending', 'processing', 'done', 'failed']:
count = conn.execute("SELECT COUNT(*) FROM queue WHERE status=?", (status,)).fetchone()[0]
stats[status] = count
conn.close()
return stats
if __name__ == "__main__":
if len(sys.argv) < 2:
print("用法: signal_queue.py enqueue '信号原文'")
print(" signal_queue.py list")
print(" signal_queue.py retry")
sys.exit(1)
cmd = sys.argv[1]
if cmd == "enqueue":
raw = sys.argv[2] if len(sys.argv) > 2 else sys.stdin.read()
row_id = enqueue(raw)
print(f"✅ 已入队 #{row_id}")
elif cmd == "list":
pending = get_pending()
if not pending:
print("队列为空,无待处理信号")
else:
for row_id, raw, retries in pending:
print(f"#{row_id} (重试{retries}次): {raw[:80]}...")
elif cmd == "retry":
"""重试所有失败信号"""
pending = get_pending()
print(f"待处理: {len(pending)}条")
for row_id, raw, retries in pending:
print(f"\n重试 #{row_id}...")
mark_processing(row_id)
import subprocess
try:
r = subprocess.run(
["python3", os.path.expanduser("~/.hermes/skills/trading/okx-auto-position/scripts/process_signal.py"), raw],
capture_output=True, text=True, timeout=120
)
if r.returncode == 0:
mark_done(row_id, r.stdout[:200])
print(f" ✅ 成功")
else:
mark_failed(row_id, r.stderr[:200])
print(f" ❌ 失败: {r.stderr[:100]}")
except Exception as e:
mark_failed(row_id, str(e))
print(f" ❌ 异常: {e}")
elif cmd == "stats":
stats = get_stats()
print(f"待处理: {stats['pending']} | 处理中: {stats['processing']} | 完成: {stats['done']} | 失败: {stats['failed']}")