#!/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']}")