Spaces:
Paused
Paused
zt p commited on
Commit ·
3be8894
1
Parent(s): b04609c
Persist monitor state across restarts
Browse files- .env.example +8 -0
- README.md +12 -1
- app/core/alerts.py +6 -1
- app/core/dashboard.py +23 -4
- app/core/database.py +15 -0
- app/core/platform.py +50 -7
- app/core/state.py +142 -0
- app/monitors/base.py +2 -1
- app/monitors/crypto_price.py +20 -8
- app/monitors/market_move.py +6 -2
- app/monitors/strategy_treasury.py +2 -2
- coinpush.py +74 -38
- requirements.txt +1 -0
.env.example
ADDED
|
@@ -0,0 +1,8 @@
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
+
# Hugging Face deployment / persistent state
|
| 2 |
+
HF_TOKEN=
|
| 3 |
+
STATE_HUB_REPO_ID=
|
| 4 |
+
STATE_HUB_REPO_TYPE=dataset
|
| 5 |
+
STATE_HUB_FILENAME=monitor_state.json
|
| 6 |
+
|
| 7 |
+
# 企业微信机器人 Webhook
|
| 8 |
+
WEWORK_BOT_WEBHOOK=
|
README.md
CHANGED
|
@@ -66,6 +66,17 @@ uvicorn app.main:app --host 0.0.0.0 --port 7860
|
|
| 66 |
|
| 67 |
如果 Space 没有挂载持久化存储,重启后 `/data` 数据可能丢失。建议在 Space 设置里挂载 Storage 后再长期运行。
|
| 68 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 69 |
## 本地运行
|
| 70 |
|
| 71 |
```bash
|
|
@@ -126,7 +137,7 @@ DISABLE_WORKER=true .venv/bin/uvicorn app.main:app --host 0.0.0.0 --port 7860
|
|
| 126 |
- `BINANCE_API_KEY` / `BINANCE_API_SECRET`:币安私有接口凭证,可选。
|
| 127 |
- `BINANCE_SPOT_BASE` / `BINANCE_FAPI_BASE` / `BINANCE_FSTREAM_BASE`:币安接口基址,可选。
|
| 128 |
- `BINANCE_ALPHA_BASE`:Binance Alpha 公共行情基址,默认 `https://www.binance.com`。
|
| 129 |
-
-
|
| 130 |
- `BITGET_API_KEY` / `BITGET_API_SECRET` / `BITGET_API_PASSPHRASE`:Bitget 凭证,可选。
|
| 131 |
- `OKX_API_KEY` / `OKX_API_SECRET` / `OKX_API_PASSPHRASE`:OKX 凭证,可选。
|
| 132 |
- `BITGET_BASE` / `OKX_BASE`:交易所接口基址,可选。
|
|
|
|
| 66 |
|
| 67 |
如果 Space 没有挂载持久化存储,重启后 `/data` 数据可能丢失。建议在 Space 设置里挂载 Storage 后再长期运行。
|
| 68 |
|
| 69 |
+
如果没有挂载 Storage,也可以配置私有状态数据集,让关键配置跨重启保留:
|
| 70 |
+
|
| 71 |
+
```text
|
| 72 |
+
HF_TOKEN=<具备写入权限的 Hugging Face token>
|
| 73 |
+
STATE_HUB_REPO_ID=<用户名>/coinpush-state
|
| 74 |
+
STATE_HUB_REPO_TYPE=dataset
|
| 75 |
+
STATE_HUB_FILENAME=monitor_state.json
|
| 76 |
+
```
|
| 77 |
+
|
| 78 |
+
状态文件会保存币价配置、全市场涨跌配置、暂停的监控项 ID 和告警冷却时间。写入状态失败时不影响监控运行,本地文件仍会保留。
|
| 79 |
+
|
| 80 |
## 本地运行
|
| 81 |
|
| 82 |
```bash
|
|
|
|
| 137 |
- `BINANCE_API_KEY` / `BINANCE_API_SECRET`:币安私有接口凭证,可选。
|
| 138 |
- `BINANCE_SPOT_BASE` / `BINANCE_FAPI_BASE` / `BINANCE_FSTREAM_BASE`:币安接口基址,可选。
|
| 139 |
- `BINANCE_ALPHA_BASE`:Binance Alpha 公共行情基址,默认 `https://www.binance.com`。
|
| 140 |
+
- 合约指标和强平流由币种配置里的 `enable_futures` 控制;币安合约接口受限时字段会降级为空。
|
| 141 |
- `BITGET_API_KEY` / `BITGET_API_SECRET` / `BITGET_API_PASSPHRASE`:Bitget 凭证,可选。
|
| 142 |
- `OKX_API_KEY` / `OKX_API_SECRET` / `OKX_API_PASSPHRASE`:OKX 凭证,可选。
|
| 143 |
- `BITGET_BASE` / `OKX_BASE`:交易所接口基址,可选。
|
app/core/alerts.py
CHANGED
|
@@ -4,13 +4,16 @@ import time
|
|
| 4 |
class SQLiteAlertManager:
|
| 5 |
"""Alert sender with cooldown and persistent alert history."""
|
| 6 |
|
| 7 |
-
def __init__(self, pusher, store, logger, stats, cooldowns=None):
|
| 8 |
self.pusher = pusher
|
| 9 |
self.store = store
|
| 10 |
self.logger = logger
|
| 11 |
self.stats = stats
|
| 12 |
self.cooldowns = cooldowns or {}
|
| 13 |
self.last_sent = {}
|
|
|
|
|
|
|
|
|
|
| 14 |
|
| 15 |
def emit(self, alerts):
|
| 16 |
alerts = sorted(alerts, key=lambda a: (a.get("direction") != "跌",))
|
|
@@ -37,6 +40,8 @@ class SQLiteAlertManager:
|
|
| 37 |
delivered=delivered,
|
| 38 |
error=None if delivered else "notifier returned false",
|
| 39 |
)
|
|
|
|
|
|
|
| 40 |
self.logger(
|
| 41 |
f"{alert.get('sev_emoji', '')} 推送[{alert.get('category')}] "
|
| 42 |
f"{alert.get('title', key)}",
|
|
|
|
| 4 |
class SQLiteAlertManager:
|
| 5 |
"""Alert sender with cooldown and persistent alert history."""
|
| 6 |
|
| 7 |
+
def __init__(self, pusher, store, logger, stats, cooldowns=None, state_store=None):
|
| 8 |
self.pusher = pusher
|
| 9 |
self.store = store
|
| 10 |
self.logger = logger
|
| 11 |
self.stats = stats
|
| 12 |
self.cooldowns = cooldowns or {}
|
| 13 |
self.last_sent = {}
|
| 14 |
+
self.state_store = state_store
|
| 15 |
+
if self.state_store is not None:
|
| 16 |
+
self.last_sent = dict(self.state_store.get("alert_sent_at") or {})
|
| 17 |
|
| 18 |
def emit(self, alerts):
|
| 19 |
alerts = sorted(alerts, key=lambda a: (a.get("direction") != "跌",))
|
|
|
|
| 40 |
delivered=delivered,
|
| 41 |
error=None if delivered else "notifier returned false",
|
| 42 |
)
|
| 43 |
+
if self.state_store is not None:
|
| 44 |
+
self.state_store.set_kv("alert_sent_at", dict(self.last_sent))
|
| 45 |
self.logger(
|
| 46 |
f"{alert.get('sev_emoji', '')} 推送[{alert.get('category')}] "
|
| 47 |
f"{alert.get('title', key)}",
|
app/core/dashboard.py
CHANGED
|
@@ -393,6 +393,11 @@ def render_dashboard():
|
|
| 393 |
overflow: auto;
|
| 394 |
overflow-wrap: anywhere;
|
| 395 |
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 396 |
.log-snippet {
|
| 397 |
max-height: 120px;
|
| 398 |
overflow: auto;
|
|
@@ -1167,10 +1172,16 @@ def render_dashboard():
|
|
| 1167 |
return chip(label, formatter ? formatter(value) : value);
|
| 1168 |
}
|
| 1169 |
|
| 1170 |
-
function rawDetails(payload, raw) {
|
| 1171 |
if (!payload && !raw) return "";
|
| 1172 |
const text = payload ? JSON.stringify(payload, null, 2) : String(raw);
|
| 1173 |
-
return `<details><summary>
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1174 |
}
|
| 1175 |
|
| 1176 |
function firstScalarChips(payload) {
|
|
@@ -1246,18 +1257,26 @@ def render_dashboard():
|
|
| 1246 |
function renderMonitorValue(monitor) {
|
| 1247 |
const payload = parseJsonField(monitor.last_value_json);
|
| 1248 |
const raw = monitor.last_value_json;
|
| 1249 |
-
const
|
|
|
|
|
|
|
| 1250 |
|
| 1251 |
if (monitor.last_error) {
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1252 |
return (
|
| 1253 |
`<div class="value-card">` +
|
| 1254 |
`<div class="error-box">${escapeHtml(monitor.last_error)}</div>` +
|
|
|
|
| 1255 |
details +
|
| 1256 |
`</div>`
|
| 1257 |
);
|
| 1258 |
}
|
| 1259 |
|
| 1260 |
-
const
|
|
|
|
| 1261 |
if (!chips.length && !details) return "<span class='muted'>暂无数据</span>";
|
| 1262 |
return (
|
| 1263 |
`<div class="value-card">` +
|
|
|
|
| 393 |
overflow: auto;
|
| 394 |
overflow-wrap: anywhere;
|
| 395 |
}
|
| 396 |
+
.stale-note {
|
| 397 |
+
color: var(--muted);
|
| 398 |
+
font-size: 12px;
|
| 399 |
+
line-height: 1.4;
|
| 400 |
+
}
|
| 401 |
.log-snippet {
|
| 402 |
max-height: 120px;
|
| 403 |
overflow: auto;
|
|
|
|
| 1172 |
return chip(label, formatter ? formatter(value) : value);
|
| 1173 |
}
|
| 1174 |
|
| 1175 |
+
function rawDetails(payload, raw, summary = "原始数据") {
|
| 1176 |
if (!payload && !raw) return "";
|
| 1177 |
const text = payload ? JSON.stringify(payload, null, 2) : String(raw);
|
| 1178 |
+
return `<details><summary>${escapeHtml(summary)}</summary><pre class="value-json">${escapeHtml(text)}</pre></details>`;
|
| 1179 |
+
}
|
| 1180 |
+
|
| 1181 |
+
function isCalibrationOnlyPayload(payload) {
|
| 1182 |
+
if (!payload || typeof payload !== "object" || Array.isArray(payload)) return false;
|
| 1183 |
+
const keys = Object.keys(payload);
|
| 1184 |
+
return keys.length === 1 && payload.calibrated === true;
|
| 1185 |
}
|
| 1186 |
|
| 1187 |
function firstScalarChips(payload) {
|
|
|
|
| 1257 |
function renderMonitorValue(monitor) {
|
| 1258 |
const payload = parseJsonField(monitor.last_value_json);
|
| 1259 |
const raw = monitor.last_value_json;
|
| 1260 |
+
const hasMarketPayload = !isCalibrationOnlyPayload(payload);
|
| 1261 |
+
const displayPayload = hasMarketPayload ? payload : null;
|
| 1262 |
+
const displayRaw = hasMarketPayload ? raw : null;
|
| 1263 |
|
| 1264 |
if (monitor.last_error) {
|
| 1265 |
+
const staleNote = displayRaw && monitor.last_success_at
|
| 1266 |
+
? `<div class="stale-note">下方为上次成功数据:${escapeHtml(formatDateTime(monitor.last_success_at))}</div>`
|
| 1267 |
+
: "";
|
| 1268 |
+
const details = rawDetails(displayPayload, displayRaw, "上次成功数据");
|
| 1269 |
return (
|
| 1270 |
`<div class="value-card">` +
|
| 1271 |
`<div class="error-box">${escapeHtml(monitor.last_error)}</div>` +
|
| 1272 |
+
staleNote +
|
| 1273 |
details +
|
| 1274 |
`</div>`
|
| 1275 |
);
|
| 1276 |
}
|
| 1277 |
|
| 1278 |
+
const details = rawDetails(displayPayload, displayRaw);
|
| 1279 |
+
const chips = monitorSummaryChips(monitor, displayPayload);
|
| 1280 |
if (!chips.length && !details) return "<span class='muted'>暂无数据</span>";
|
| 1281 |
return (
|
| 1282 |
`<div class="value-card">` +
|
app/core/database.py
CHANGED
|
@@ -125,6 +125,21 @@ class MonitorStore:
|
|
| 125 |
),
|
| 126 |
)
|
| 127 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 128 |
def get_monitor(self, monitor_id):
|
| 129 |
with self._lock:
|
| 130 |
row = self._conn.execute(
|
|
|
|
| 125 |
),
|
| 126 |
)
|
| 127 |
|
| 128 |
+
def apply_paused_ids(self, monitor_ids):
|
| 129 |
+
ids = sorted(set(monitor_ids or []))
|
| 130 |
+
if not ids:
|
| 131 |
+
return
|
| 132 |
+
now = utc_now_iso()
|
| 133 |
+
with self._lock, self._conn:
|
| 134 |
+
self._conn.executemany(
|
| 135 |
+
"""
|
| 136 |
+
UPDATE monitors
|
| 137 |
+
SET paused=1, status='paused', next_run_at=NULL, updated_at=?
|
| 138 |
+
WHERE id=?
|
| 139 |
+
""",
|
| 140 |
+
[(now, monitor_id) for monitor_id in ids],
|
| 141 |
+
)
|
| 142 |
+
|
| 143 |
def get_monitor(self, monitor_id):
|
| 144 |
with self._lock:
|
| 145 |
row = self._conn.execute(
|
app/core/platform.py
CHANGED
|
@@ -1,7 +1,9 @@
|
|
|
|
|
| 1 |
import os
|
| 2 |
import threading
|
| 3 |
|
| 4 |
from app.core.database import MonitorStore
|
|
|
|
| 5 |
from app.monitors.crypto_price import CryptoPriceMonitor
|
| 6 |
from app.monitors.market_move import MarketMoveMonitor
|
| 7 |
from app.monitors.strategy_treasury import StrategyTreasuryMonitor
|
|
@@ -15,18 +17,44 @@ class MonitoringPlatform:
|
|
| 15 |
if not os.path.isdir(data_dir) or not os.access(data_dir, os.W_OK):
|
| 16 |
data_dir = os.path.join(os.getcwd(), "data")
|
| 17 |
self.data_dir = data_dir
|
|
|
|
| 18 |
self.store = MonitorStore(os.getenv("MONITOR_DB_PATH", os.path.join(data_dir, "monitor.db")))
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 19 |
self._lock = threading.RLock()
|
| 20 |
-
self.crypto = CryptoPriceMonitor(self.store)
|
| 21 |
-
self.strategy_treasury = StrategyTreasuryMonitor(self.store)
|
| 22 |
-
self.market_move = MarketMoveMonitor(self.store)
|
| 23 |
self.plugins = {
|
| 24 |
"crypto_price": self.crypto,
|
| 25 |
"strategy_treasury": self.strategy_treasury,
|
| 26 |
"market_move": self.market_move,
|
| 27 |
}
|
|
|
|
|
|
|
|
|
|
| 28 |
self.started = False
|
| 29 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 30 |
def start(self):
|
| 31 |
with self._lock:
|
| 32 |
if self.started:
|
|
@@ -46,6 +74,11 @@ class MonitoringPlatform:
|
|
| 46 |
return {
|
| 47 |
"started": self.started,
|
| 48 |
"data_dir": self.data_dir,
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 49 |
"monitors": self.store.list_monitors(),
|
| 50 |
"alerts": self.store.recent_alerts(20),
|
| 51 |
"events": self.store.recent_events(40),
|
|
@@ -71,14 +104,22 @@ class MonitoringPlatform:
|
|
| 71 |
}
|
| 72 |
|
| 73 |
def set_monitor_paused(self, monitor_id, paused):
|
| 74 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 75 |
|
| 76 |
def get_config(self):
|
| 77 |
return self.crypto.get_config()
|
| 78 |
|
| 79 |
def update_config(self, cfg):
|
| 80 |
-
self.crypto.update_config(cfg)
|
| 81 |
-
|
|
|
|
| 82 |
|
| 83 |
def calibrate_all(self):
|
| 84 |
self.crypto.calibrate_all()
|
|
@@ -91,4 +132,6 @@ class MonitoringPlatform:
|
|
| 91 |
return self.market_move.get_config()
|
| 92 |
|
| 93 |
def update_market_move_config(self, cfg):
|
| 94 |
-
|
|
|
|
|
|
|
|
|
| 1 |
+
import json
|
| 2 |
import os
|
| 3 |
import threading
|
| 4 |
|
| 5 |
from app.core.database import MonitorStore
|
| 6 |
+
from app.core.state import HubStateStore
|
| 7 |
from app.monitors.crypto_price import CryptoPriceMonitor
|
| 8 |
from app.monitors.market_move import MarketMoveMonitor
|
| 9 |
from app.monitors.strategy_treasury import StrategyTreasuryMonitor
|
|
|
|
| 17 |
if not os.path.isdir(data_dir) or not os.access(data_dir, os.W_OK):
|
| 18 |
data_dir = os.path.join(os.getcwd(), "data")
|
| 19 |
self.data_dir = data_dir
|
| 20 |
+
self.state_store = HubStateStore(data_dir)
|
| 21 |
self.store = MonitorStore(os.getenv("MONITOR_DB_PATH", os.path.join(data_dir, "monitor.db")))
|
| 22 |
+
crypto_config = self.state_store.get("crypto_config")
|
| 23 |
+
if crypto_config is not None:
|
| 24 |
+
self._write_crypto_config(crypto_config)
|
| 25 |
+
self.store.set_kv("crypto_config", crypto_config)
|
| 26 |
+
market_move_config = self.state_store.get("market_move:config")
|
| 27 |
+
if market_move_config is not None:
|
| 28 |
+
self.store.set_kv("market_move:config", market_move_config)
|
| 29 |
+
market_move_alert_state = self.state_store.get("market_move:alert_state")
|
| 30 |
+
if market_move_alert_state is not None:
|
| 31 |
+
self.store.set_kv("market_move:alert_state", market_move_alert_state)
|
| 32 |
self._lock = threading.RLock()
|
| 33 |
+
self.crypto = CryptoPriceMonitor(self.store, self.state_store)
|
| 34 |
+
self.strategy_treasury = StrategyTreasuryMonitor(self.store, self.state_store)
|
| 35 |
+
self.market_move = MarketMoveMonitor(self.store, self.state_store)
|
| 36 |
self.plugins = {
|
| 37 |
"crypto_price": self.crypto,
|
| 38 |
"strategy_treasury": self.strategy_treasury,
|
| 39 |
"market_move": self.market_move,
|
| 40 |
}
|
| 41 |
+
paused_ids = self.state_store.get("paused_monitor_ids") or []
|
| 42 |
+
if paused_ids:
|
| 43 |
+
self.store.apply_paused_ids(paused_ids)
|
| 44 |
self.started = False
|
| 45 |
|
| 46 |
+
def _write_crypto_config(self, cfg):
|
| 47 |
+
try:
|
| 48 |
+
from coinpush import CONFIG_FILE
|
| 49 |
+
|
| 50 |
+
os.makedirs(os.path.dirname(CONFIG_FILE), exist_ok=True)
|
| 51 |
+
temp_path = CONFIG_FILE + ".tmp"
|
| 52 |
+
with open(temp_path, "w", encoding="utf-8") as handle:
|
| 53 |
+
json.dump(cfg, handle, ensure_ascii=False, indent=2)
|
| 54 |
+
os.replace(temp_path, CONFIG_FILE)
|
| 55 |
+
except Exception as exc:
|
| 56 |
+
self.state_store.last_error = str(exc)
|
| 57 |
+
|
| 58 |
def start(self):
|
| 59 |
with self._lock:
|
| 60 |
if self.started:
|
|
|
|
| 74 |
return {
|
| 75 |
"started": self.started,
|
| 76 |
"data_dir": self.data_dir,
|
| 77 |
+
"state_sync": {
|
| 78 |
+
"enabled": self.state_store.enabled,
|
| 79 |
+
"remote_synced": self.state_store.remote_synced,
|
| 80 |
+
"last_error": self.state_store.last_error,
|
| 81 |
+
},
|
| 82 |
"monitors": self.store.list_monitors(),
|
| 83 |
"alerts": self.store.recent_alerts(20),
|
| 84 |
"events": self.store.recent_events(40),
|
|
|
|
| 104 |
}
|
| 105 |
|
| 106 |
def set_monitor_paused(self, monitor_id, paused):
|
| 107 |
+
result = self.store.set_monitor_paused(monitor_id, paused)
|
| 108 |
+
paused_ids = set(self.state_store.get("paused_monitor_ids") or [])
|
| 109 |
+
if paused:
|
| 110 |
+
paused_ids.add(monitor_id)
|
| 111 |
+
else:
|
| 112 |
+
paused_ids.discard(monitor_id)
|
| 113 |
+
self.state_store.update(paused_monitor_ids=sorted(paused_ids))
|
| 114 |
+
return result
|
| 115 |
|
| 116 |
def get_config(self):
|
| 117 |
return self.crypto.get_config()
|
| 118 |
|
| 119 |
def update_config(self, cfg):
|
| 120 |
+
result = self.crypto.update_config(cfg)
|
| 121 |
+
self.state_store.set_kv("crypto_config", result)
|
| 122 |
+
return result
|
| 123 |
|
| 124 |
def calibrate_all(self):
|
| 125 |
self.crypto.calibrate_all()
|
|
|
|
| 132 |
return self.market_move.get_config()
|
| 133 |
|
| 134 |
def update_market_move_config(self, cfg):
|
| 135 |
+
result = self.market_move.update_config(cfg)
|
| 136 |
+
self.state_store.set_kv("market_move:config", result)
|
| 137 |
+
return result
|
app/core/state.py
ADDED
|
@@ -0,0 +1,142 @@
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
+
import copy
|
| 2 |
+
import io
|
| 3 |
+
import json
|
| 4 |
+
import os
|
| 5 |
+
import tempfile
|
| 6 |
+
import threading
|
| 7 |
+
from datetime import datetime, timezone
|
| 8 |
+
|
| 9 |
+
|
| 10 |
+
def utc_now_iso():
|
| 11 |
+
return datetime.now(timezone.utc).isoformat(timespec="seconds")
|
| 12 |
+
|
| 13 |
+
|
| 14 |
+
class HubStateStore:
|
| 15 |
+
"""Small JSON state store backed by a private Hugging Face dataset."""
|
| 16 |
+
|
| 17 |
+
def __init__(self, data_dir):
|
| 18 |
+
self.repo_id = os.getenv("STATE_HUB_REPO_ID")
|
| 19 |
+
self.repo_type = os.getenv("STATE_HUB_REPO_TYPE", "dataset")
|
| 20 |
+
self.filename = os.getenv("STATE_HUB_FILENAME", "monitor_state.json")
|
| 21 |
+
token = os.getenv("HF_TOKEN")
|
| 22 |
+
self.path = os.path.join(data_dir, self.filename)
|
| 23 |
+
self._lock = threading.RLock()
|
| 24 |
+
self._api = None
|
| 25 |
+
self.last_error = None
|
| 26 |
+
self.remote_synced = False
|
| 27 |
+
|
| 28 |
+
if self.repo_id and token:
|
| 29 |
+
try:
|
| 30 |
+
from huggingface_hub import HfApi
|
| 31 |
+
|
| 32 |
+
self._api = HfApi(token=token)
|
| 33 |
+
except Exception as exc:
|
| 34 |
+
self.last_error = str(exc)
|
| 35 |
+
|
| 36 |
+
self._state = self._load()
|
| 37 |
+
|
| 38 |
+
@property
|
| 39 |
+
def enabled(self):
|
| 40 |
+
return self._api is not None and bool(self.repo_id)
|
| 41 |
+
|
| 42 |
+
def _empty(self):
|
| 43 |
+
return {"updated_at": "1970-01-01T00:00:00+00:00", "paused_monitor_ids": [], "values": {}}
|
| 44 |
+
|
| 45 |
+
def _decode(self, path):
|
| 46 |
+
with open(path, "r", encoding="utf-8") as handle:
|
| 47 |
+
state = json.load(handle)
|
| 48 |
+
if not isinstance(state, dict):
|
| 49 |
+
raise ValueError("state root must be an object")
|
| 50 |
+
state.setdefault("paused_monitor_ids", [])
|
| 51 |
+
state.setdefault("values", {})
|
| 52 |
+
if not isinstance(state["paused_monitor_ids"], list):
|
| 53 |
+
state["paused_monitor_ids"] = []
|
| 54 |
+
if not isinstance(state["values"], dict):
|
| 55 |
+
state["values"] = {}
|
| 56 |
+
return state
|
| 57 |
+
|
| 58 |
+
def _write_local(self, state):
|
| 59 |
+
os.makedirs(os.path.dirname(self.path), exist_ok=True)
|
| 60 |
+
fd, temp_path = tempfile.mkstemp(prefix=".monitor-state-", dir=os.path.dirname(self.path))
|
| 61 |
+
try:
|
| 62 |
+
with os.fdopen(fd, "w", encoding="utf-8") as handle:
|
| 63 |
+
json.dump(state, handle, ensure_ascii=False, indent=2, sort_keys=True)
|
| 64 |
+
os.replace(temp_path, self.path)
|
| 65 |
+
except Exception:
|
| 66 |
+
try:
|
| 67 |
+
os.unlink(temp_path)
|
| 68 |
+
except OSError:
|
| 69 |
+
pass
|
| 70 |
+
raise
|
| 71 |
+
|
| 72 |
+
def _load(self):
|
| 73 |
+
local = self._empty()
|
| 74 |
+
if os.path.exists(self.path):
|
| 75 |
+
try:
|
| 76 |
+
local = self._decode(self.path)
|
| 77 |
+
except Exception as exc:
|
| 78 |
+
self.last_error = str(exc)
|
| 79 |
+
|
| 80 |
+
if not self.enabled:
|
| 81 |
+
return local
|
| 82 |
+
|
| 83 |
+
try:
|
| 84 |
+
from huggingface_hub import hf_hub_download
|
| 85 |
+
|
| 86 |
+
remote_path = hf_hub_download(
|
| 87 |
+
repo_id=self.repo_id,
|
| 88 |
+
repo_type=self.repo_type,
|
| 89 |
+
filename=self.filename,
|
| 90 |
+
force_download=True,
|
| 91 |
+
)
|
| 92 |
+
remote = self._decode(remote_path)
|
| 93 |
+
if str(remote.get("updated_at") or "") >= str(local.get("updated_at") or ""):
|
| 94 |
+
self.remote_synced = True
|
| 95 |
+
self._write_local(remote)
|
| 96 |
+
return remote
|
| 97 |
+
return local
|
| 98 |
+
except Exception as exc:
|
| 99 |
+
self.last_error = str(exc)
|
| 100 |
+
return local
|
| 101 |
+
|
| 102 |
+
def _upload(self, state):
|
| 103 |
+
if not self.enabled:
|
| 104 |
+
return False
|
| 105 |
+
try:
|
| 106 |
+
payload = json.dumps(state, ensure_ascii=False, indent=2, sort_keys=True).encode("utf-8")
|
| 107 |
+
self._api.upload_file(
|
| 108 |
+
path_or_fileobj=io.BytesIO(payload),
|
| 109 |
+
path_in_repo=self.filename,
|
| 110 |
+
repo_id=self.repo_id,
|
| 111 |
+
repo_type=self.repo_type,
|
| 112 |
+
commit_message="Update monitor state",
|
| 113 |
+
)
|
| 114 |
+
self.last_error = None
|
| 115 |
+
self.remote_synced = True
|
| 116 |
+
return True
|
| 117 |
+
except Exception as exc:
|
| 118 |
+
self.last_error = str(exc)
|
| 119 |
+
self.remote_synced = False
|
| 120 |
+
return False
|
| 121 |
+
|
| 122 |
+
def get(self, key, default=None):
|
| 123 |
+
with self._lock:
|
| 124 |
+
if key == "paused_monitor_ids":
|
| 125 |
+
return copy.deepcopy(self._state.get(key, []))
|
| 126 |
+
return copy.deepcopy(self._state.get("values", {}).get(key, default))
|
| 127 |
+
|
| 128 |
+
def update(self, **updates):
|
| 129 |
+
with self._lock:
|
| 130 |
+
self._state.update(updates)
|
| 131 |
+
self._state["updated_at"] = utc_now_iso()
|
| 132 |
+
state = copy.deepcopy(self._state)
|
| 133 |
+
self._write_local(state)
|
| 134 |
+
self._upload(state)
|
| 135 |
+
|
| 136 |
+
def set_kv(self, key, value):
|
| 137 |
+
with self._lock:
|
| 138 |
+
self._state.setdefault("values", {})[key] = copy.deepcopy(value)
|
| 139 |
+
self._state["updated_at"] = utc_now_iso()
|
| 140 |
+
state = copy.deepcopy(self._state)
|
| 141 |
+
self._write_local(state)
|
| 142 |
+
self._upload(state)
|
app/monitors/base.py
CHANGED
|
@@ -3,8 +3,9 @@ class BaseMonitor:
|
|
| 3 |
|
| 4 |
monitor_type = "base"
|
| 5 |
|
| 6 |
-
def __init__(self, store):
|
| 7 |
self.store = store
|
|
|
|
| 8 |
|
| 9 |
def start(self):
|
| 10 |
raise NotImplementedError
|
|
|
|
| 3 |
|
| 4 |
monitor_type = "base"
|
| 5 |
|
| 6 |
+
def __init__(self, store, state_store=None):
|
| 7 |
self.store = store
|
| 8 |
+
self.state_store = state_store
|
| 9 |
|
| 10 |
def start(self):
|
| 11 |
raise NotImplementedError
|
app/monitors/crypto_price.py
CHANGED
|
@@ -29,8 +29,8 @@ class CryptoPriceMonitor(BaseMonitor):
|
|
| 29 |
|
| 30 |
monitor_type = "crypto_price"
|
| 31 |
|
| 32 |
-
def __init__(self, store):
|
| 33 |
-
super().__init__(store)
|
| 34 |
self.logs = deque(maxlen=200)
|
| 35 |
self.status_text = "初始化中..."
|
| 36 |
self.next_wakeup = None
|
|
@@ -59,6 +59,14 @@ class CryptoPriceMonitor(BaseMonitor):
|
|
| 59 |
self.log,
|
| 60 |
)
|
| 61 |
self.okx_adapter = legacy.OKXMarketAdapter(self.okx, self.log)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 62 |
self.binance_futures_adapter = legacy.BinanceFuturesMarketAdapter(self.api)
|
| 63 |
self.binance_futures_with_fallback = legacy.FallbackMarketAdapter(
|
| 64 |
self.binance_futures_adapter,
|
|
@@ -80,6 +88,7 @@ class CryptoPriceMonitor(BaseMonitor):
|
|
| 80 |
self.log,
|
| 81 |
self.stats,
|
| 82 |
cooldowns=legacy.COOLDOWN_MIN,
|
|
|
|
| 83 |
)
|
| 84 |
self.liq = legacy.LiquidationTracker(self._futures_symbols(), self.log)
|
| 85 |
self.cross_ex = legacy.CrossExchangeMonitor(self.api, self.bitget, self.okx, self.log)
|
|
@@ -161,7 +170,7 @@ class CryptoPriceMonitor(BaseMonitor):
|
|
| 161 |
return self.alpha_with_futures
|
| 162 |
if source == "okx":
|
| 163 |
return self.okx_adapter
|
| 164 |
-
return self.
|
| 165 |
|
| 166 |
def _futures_symbols(self):
|
| 167 |
symbols = []
|
|
@@ -169,9 +178,7 @@ class CryptoPriceMonitor(BaseMonitor):
|
|
| 169 |
if not coin.get("enable_futures") or not coin.get("futures_symbol"):
|
| 170 |
continue
|
| 171 |
source = coin.get("data_source", "binance")
|
| 172 |
-
if source in ("binance_futures", "binance_alpha")
|
| 173 |
-
source == "binance" and legacy.FUTURES_ENABLED
|
| 174 |
-
):
|
| 175 |
symbols.append(coin["futures_symbol"])
|
| 176 |
return symbols
|
| 177 |
|
|
@@ -252,7 +259,7 @@ class CryptoPriceMonitor(BaseMonitor):
|
|
| 252 |
self.log(f"开始校准 {name}", monitor_id=coin_monitor_id(name))
|
| 253 |
overrides = legacy.calibrate_thresholds(self._api_for(coin), coin, self.log, scale)
|
| 254 |
cfg["coins"][name] = legacy._deep_merge(coin, overrides)
|
| 255 |
-
self.
|
| 256 |
except Exception as exc:
|
| 257 |
self._record_failure(coin_monitor_id(name), f"{name} 校准失败: {exc}")
|
| 258 |
cfg["global"]["last_calibrate_ts"] = time.time()
|
|
@@ -480,7 +487,12 @@ class CryptoPriceMonitor(BaseMonitor):
|
|
| 480 |
)
|
| 481 |
else:
|
| 482 |
self.stats["api_fails"] += 1
|
| 483 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 484 |
total_alerts.extend(alerts)
|
| 485 |
except Exception as exc:
|
| 486 |
self._record_failure(monitor_id, f"{name} 评估异常: {exc}")
|
|
|
|
| 29 |
|
| 30 |
monitor_type = "crypto_price"
|
| 31 |
|
| 32 |
+
def __init__(self, store, state_store=None):
|
| 33 |
+
super().__init__(store, state_store)
|
| 34 |
self.logs = deque(maxlen=200)
|
| 35 |
self.status_text = "初始化中..."
|
| 36 |
self.next_wakeup = None
|
|
|
|
| 59 |
self.log,
|
| 60 |
)
|
| 61 |
self.okx_adapter = legacy.OKXMarketAdapter(self.okx, self.log)
|
| 62 |
+
self.binance_spot_with_fallback = legacy.FallbackMarketAdapter(
|
| 63 |
+
self.api,
|
| 64 |
+
self.okx_adapter,
|
| 65 |
+
self.log,
|
| 66 |
+
primary_label="Binance现货主源",
|
| 67 |
+
fallback_label="OKX",
|
| 68 |
+
fallback_symbol_resolver=legacy.okx_fallback_symbol,
|
| 69 |
+
)
|
| 70 |
self.binance_futures_adapter = legacy.BinanceFuturesMarketAdapter(self.api)
|
| 71 |
self.binance_futures_with_fallback = legacy.FallbackMarketAdapter(
|
| 72 |
self.binance_futures_adapter,
|
|
|
|
| 88 |
self.log,
|
| 89 |
self.stats,
|
| 90 |
cooldowns=legacy.COOLDOWN_MIN,
|
| 91 |
+
state_store=state_store,
|
| 92 |
)
|
| 93 |
self.liq = legacy.LiquidationTracker(self._futures_symbols(), self.log)
|
| 94 |
self.cross_ex = legacy.CrossExchangeMonitor(self.api, self.bitget, self.okx, self.log)
|
|
|
|
| 170 |
return self.alpha_with_futures
|
| 171 |
if source == "okx":
|
| 172 |
return self.okx_adapter
|
| 173 |
+
return self.binance_spot_with_fallback
|
| 174 |
|
| 175 |
def _futures_symbols(self):
|
| 176 |
symbols = []
|
|
|
|
| 178 |
if not coin.get("enable_futures") or not coin.get("futures_symbol"):
|
| 179 |
continue
|
| 180 |
source = coin.get("data_source", "binance")
|
| 181 |
+
if source in ("binance", "binance_futures", "binance_alpha"):
|
|
|
|
|
|
|
| 182 |
symbols.append(coin["futures_symbol"])
|
| 183 |
return symbols
|
| 184 |
|
|
|
|
| 259 |
self.log(f"开始校准 {name}", monitor_id=coin_monitor_id(name))
|
| 260 |
overrides = legacy.calibrate_thresholds(self._api_for(coin), coin, self.log, scale)
|
| 261 |
cfg["coins"][name] = legacy._deep_merge(coin, overrides)
|
| 262 |
+
self.store.add_event(coin_monitor_id(name), "info", "校准完成")
|
| 263 |
except Exception as exc:
|
| 264 |
self._record_failure(coin_monitor_id(name), f"{name} 校准失败: {exc}")
|
| 265 |
cfg["global"]["last_calibrate_ts"] = time.time()
|
|
|
|
| 487 |
)
|
| 488 |
else:
|
| 489 |
self.stats["api_fails"] += 1
|
| 490 |
+
data_error = getattr(evaluator, "last_data_error", None)
|
| 491 |
+
error = (
|
| 492 |
+
f"数据不足或接口无返回:{data_error}"
|
| 493 |
+
if data_error else "数据不足或接口无返回"
|
| 494 |
+
)
|
| 495 |
+
self._record_failure(monitor_id, error)
|
| 496 |
total_alerts.extend(alerts)
|
| 497 |
except Exception as exc:
|
| 498 |
self._record_failure(monitor_id, f"{name} 评估异常: {exc}")
|
app/monitors/market_move.py
CHANGED
|
@@ -1185,8 +1185,8 @@ class MarketMoveMonitor(BaseMonitor):
|
|
| 1185 |
|
| 1186 |
monitor_type = "market_move"
|
| 1187 |
|
| 1188 |
-
def __init__(self, store):
|
| 1189 |
-
super().__init__(store)
|
| 1190 |
self.config = self._load_config()
|
| 1191 |
self.enabled = bool(self.config.get("enabled", True))
|
| 1192 |
self.interval = int(self.config.get("interval_sec", 300))
|
|
@@ -1230,6 +1230,8 @@ class MarketMoveMonitor(BaseMonitor):
|
|
| 1230 |
self.threshold_pct = float(cfg.get("threshold_pct", 30.0))
|
| 1231 |
self.client.update_config(cfg)
|
| 1232 |
self.store.set_kv(CONFIG_KEY, cfg)
|
|
|
|
|
|
|
| 1233 |
self._register_monitors()
|
| 1234 |
self.log("全市场涨跌监控配置已更新")
|
| 1235 |
return self.get_config()
|
|
@@ -1459,6 +1461,8 @@ class MarketMoveMonitor(BaseMonitor):
|
|
| 1459 |
changed = True
|
| 1460 |
if changed:
|
| 1461 |
self.store.set_kv(ALERT_STATE_KEY, state)
|
|
|
|
|
|
|
| 1462 |
return sent
|
| 1463 |
|
| 1464 |
def _send_move_alert(self, snapshot, window, change_info, direction, rule):
|
|
|
|
| 1185 |
|
| 1186 |
monitor_type = "market_move"
|
| 1187 |
|
| 1188 |
+
def __init__(self, store, state_store=None):
|
| 1189 |
+
super().__init__(store, state_store)
|
| 1190 |
self.config = self._load_config()
|
| 1191 |
self.enabled = bool(self.config.get("enabled", True))
|
| 1192 |
self.interval = int(self.config.get("interval_sec", 300))
|
|
|
|
| 1230 |
self.threshold_pct = float(cfg.get("threshold_pct", 30.0))
|
| 1231 |
self.client.update_config(cfg)
|
| 1232 |
self.store.set_kv(CONFIG_KEY, cfg)
|
| 1233 |
+
if self.state_store is not None:
|
| 1234 |
+
self.state_store.set_kv(CONFIG_KEY, cfg)
|
| 1235 |
self._register_monitors()
|
| 1236 |
self.log("全市场涨跌监控配置已更新")
|
| 1237 |
return self.get_config()
|
|
|
|
| 1461 |
changed = True
|
| 1462 |
if changed:
|
| 1463 |
self.store.set_kv(ALERT_STATE_KEY, state)
|
| 1464 |
+
if self.state_store is not None:
|
| 1465 |
+
self.state_store.set_kv(ALERT_STATE_KEY, state)
|
| 1466 |
return sent
|
| 1467 |
|
| 1468 |
def _send_move_alert(self, snapshot, window, change_info, direction, rule):
|
app/monitors/strategy_treasury.py
CHANGED
|
@@ -245,8 +245,8 @@ class StrategyTreasuryClient:
|
|
| 245 |
class StrategyTreasuryMonitor(BaseMonitor):
|
| 246 |
monitor_type = "strategy_treasury"
|
| 247 |
|
| 248 |
-
def __init__(self, store):
|
| 249 |
-
super().__init__(store)
|
| 250 |
self.interval = int(os.getenv("STRATEGY_TREASURY_INTERVAL_SEC", "60"))
|
| 251 |
self.enabled = os.getenv("ENABLE_STRATEGY_TREASURY", "true").lower() not in ("0", "false", "no", "off")
|
| 252 |
self._stop_event = threading.Event()
|
|
|
|
| 245 |
class StrategyTreasuryMonitor(BaseMonitor):
|
| 246 |
monitor_type = "strategy_treasury"
|
| 247 |
|
| 248 |
+
def __init__(self, store, state_store=None):
|
| 249 |
+
super().__init__(store, state_store)
|
| 250 |
self.interval = int(os.getenv("STRATEGY_TREASURY_INTERVAL_SEC", "60"))
|
| 251 |
self.enabled = os.getenv("ENABLE_STRATEGY_TREASURY", "true").lower() not in ("0", "false", "no", "off")
|
| 252 |
self._stop_event = threading.Event()
|
coinpush.py
CHANGED
|
@@ -89,12 +89,6 @@ UPBIT_BASE = os.getenv("UPBIT_BASE", "https://api.upbit.com")
|
|
| 89 |
HYPERLIQUID_BASE = os.getenv("HYPERLIQUID_BASE", "https://api.hyperliquid.xyz")
|
| 90 |
ASTER_FAPI_BASE = os.getenv("ASTER_FAPI_BASE", "https://fapi.asterdex.com")
|
| 91 |
|
| 92 |
-
# 合约功能全局开关(默认关闭)。币安合约接口/强平流对 US 区域可能返回 451。
|
| 93 |
-
# 该开关控制现货主行情币种的合约指标与强平 websocket;若标的本身以 Binance
|
| 94 |
-
# 合约作为主行情源(如 STRC),会优先访问 /fapi,受限时自动使用 OKX 备用行情。
|
| 95 |
-
# 若自带可访问合约接口的代理,设置 BINANCE_FAPI_BASE 指向代理即可恢复 Binance 主源。
|
| 96 |
-
FUTURES_ENABLED = os.getenv("ENABLE_FUTURES", "false").strip().lower() in ("1", "true", "yes", "on")
|
| 97 |
-
|
| 98 |
# 配置持久化文件
|
| 99 |
CURRENT_DIR = os.path.dirname(os.path.abspath(__file__))
|
| 100 |
LOCAL_CONFIG_FILE = os.path.join(CURRENT_DIR, "crypto_config.json")
|
|
@@ -120,17 +114,15 @@ def _btc_defaults():
|
|
| 120 |
return {
|
| 121 |
"symbol": "BTCUSDT",
|
| 122 |
"futures_symbol": "BTCUSDT",
|
| 123 |
-
# 数据来源:binance(默认)或 okx。
|
| 124 |
-
# 合约信息(资金费率/持仓量)不受 FUTURES_ENABLED(仅针对币安US限制)约束。
|
| 125 |
"data_source": "binance",
|
| 126 |
# 首页核心展示:同一时间只应有一个币种开启;旧配置默认 GRAM。
|
| 127 |
"core_display": False,
|
| 128 |
# 休市属性:股票类标的(如 STRC)休市时价格长时间不变,期间不推送任何提醒
|
| 129 |
"market_close": False,
|
| 130 |
"market_close_static_secs": 600,
|
| 131 |
-
#
|
| 132 |
-
# 返回 451
|
| 133 |
-
# 不受限制。如自带可访问合约接口的代理,可在页面重新开启。
|
| 134 |
"enable_futures": False,
|
| 135 |
# 个人关键价位(0 表示未设置,不监控)
|
| 136 |
"cost_price": 0.0,
|
|
@@ -1568,7 +1560,7 @@ class FallbackMarketAdapter:
|
|
| 1568 |
|
| 1569 |
def __init__(self, primary, fallback, logger, fallback_symbols=None,
|
| 1570 |
primary_label="主源", fallback_label="备用源",
|
| 1571 |
-
allow_symbol_guess=True):
|
| 1572 |
self.primary = primary
|
| 1573 |
self.fallback = fallback
|
| 1574 |
self.logger = logger
|
|
@@ -1576,9 +1568,12 @@ class FallbackMarketAdapter:
|
|
| 1576 |
self.primary_label = primary_label
|
| 1577 |
self.fallback_label = fallback_label
|
| 1578 |
self.allow_symbol_guess = allow_symbol_guess
|
|
|
|
| 1579 |
self._fallback_logged = set()
|
| 1580 |
|
| 1581 |
-
def _fallback_symbol(self, symbol):
|
|
|
|
|
|
|
| 1582 |
if symbol in self.fallback_symbols:
|
| 1583 |
return self.fallback_symbols[symbol]
|
| 1584 |
if self.allow_symbol_guess and symbol.endswith("USDT") and len(symbol) > 4:
|
|
@@ -1593,7 +1588,7 @@ class FallbackMarketAdapter:
|
|
| 1593 |
if result:
|
| 1594 |
return result
|
| 1595 |
|
| 1596 |
-
fb_symbol = self._fallback_symbol(symbol)
|
| 1597 |
if not fb_symbol:
|
| 1598 |
return None
|
| 1599 |
fallback_fn = getattr(self.fallback, method, None)
|
|
@@ -1634,6 +1629,20 @@ class FallbackMarketAdapter:
|
|
| 1634 |
return self._call("funding_rate_history", symbol, limit)
|
| 1635 |
|
| 1636 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1637 |
# --- 7. 实时强平监控(websocket,可选)---
|
| 1638 |
class LiquidationTracker:
|
| 1639 |
"""
|
|
@@ -1876,12 +1885,10 @@ def calibrate_thresholds(api, coin_cfg, logger=print, sensitivity=1.0):
|
|
| 1876 |
out["liq_5m_usdt"] = round(max(avg5 * 0.5, 10000))
|
| 1877 |
out["liq_severe_usdt"] = round(out["liq_5m_usdt"] * 4)
|
| 1878 |
|
| 1879 |
-
# 7) 资金费率:历史费率分位(合约启用时
|
| 1880 |
-
_ds = coin_cfg.get("data_source", "binance")
|
| 1881 |
if (
|
| 1882 |
coin_cfg.get("enable_futures")
|
| 1883 |
and coin_cfg.get("futures_symbol")
|
| 1884 |
-
and (_ds in ("okx", "binance_futures", "binance_alpha") or FUTURES_ENABLED)
|
| 1885 |
):
|
| 1886 |
fr = api.funding_rate_history(coin_cfg["futures_symbol"], 500)
|
| 1887 |
if fr:
|
|
@@ -2022,6 +2029,7 @@ class CoinEvaluator:
|
|
| 2022 |
self.last_oi_changes = {}
|
| 2023 |
self.last_liq_5m_usdt = None
|
| 2024 |
self.last_rsi = None
|
|
|
|
| 2025 |
|
| 2026 |
def evaluate(self, cfg):
|
| 2027 |
sym = cfg["symbol"]
|
|
@@ -2029,18 +2037,40 @@ class CoinEvaluator:
|
|
| 2029 |
|
| 2030 |
ticker = self.api.ticker_24h(sym)
|
| 2031 |
k1 = self.api.klines(sym, "1m", 62)
|
| 2032 |
-
|
| 2033 |
-
|
| 2034 |
-
|
| 2035 |
-
|
| 2036 |
-
|
| 2037 |
-
|
| 2038 |
-
|
| 2039 |
-
|
| 2040 |
-
|
| 2041 |
-
|
| 2042 |
-
|
| 2043 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 2044 |
volume5m = sum(qvols[-5:]) if len(qvols) >= 5 else 0.0
|
| 2045 |
avg5m_volume = qvol24 / 288.0 if qvol24 > 0 else 0.0
|
| 2046 |
volume_ratio_5m = (volume5m / avg5m_volume) if avg5m_volume > 0 else None
|
|
@@ -2057,9 +2087,9 @@ class CoinEvaluator:
|
|
| 2057 |
"data_source": cfg.get("data_source", self.data_source),
|
| 2058 |
"enable_futures": bool(cfg.get("enable_futures")),
|
| 2059 |
"price": price,
|
| 2060 |
-
"high24": high24,
|
| 2061 |
-
"low24": low24,
|
| 2062 |
-
"quoteVolume24h": qvol24,
|
| 2063 |
"volume5m": volume5m,
|
| 2064 |
"volumeRatio5m": volume_ratio_5m,
|
| 2065 |
"priceChange": price_changes,
|
|
@@ -2068,6 +2098,7 @@ class CoinEvaluator:
|
|
| 2068 |
"openInterest": self.last_open_interest,
|
| 2069 |
"openInterestChange": copy.deepcopy(self.last_oi_changes),
|
| 2070 |
"liquidation5m": self.last_liq_5m_usdt,
|
|
|
|
| 2071 |
"updatedAt": datetime.now(timezone.utc).isoformat(),
|
| 2072 |
}
|
| 2073 |
|
|
@@ -2170,11 +2201,10 @@ class CoinEvaluator:
|
|
| 2170 |
# 11) RSI
|
| 2171 |
self._check_rsi(cfg, sym, head, add)
|
| 2172 |
|
| 2173 |
-
# 12) 合约类
|
| 2174 |
if (
|
| 2175 |
cfg.get("enable_futures")
|
| 2176 |
and cfg.get("futures_symbol")
|
| 2177 |
-
and (self.data_source in ("okx", "binance_futures", "binance_alpha") or FUTURES_ENABLED)
|
| 2178 |
):
|
| 2179 |
self._check_futures(cfg, head, add)
|
| 2180 |
|
|
@@ -2659,6 +2689,14 @@ class CryptoMonitor:
|
|
| 2659 |
self.bitget = BitgetAPI(self.log)
|
| 2660 |
self.okx = OKXAPI(self.log)
|
| 2661 |
self.okx_adapter = OKXMarketAdapter(self.okx, self.log)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 2662 |
self.binance_futures_adapter = BinanceFuturesMarketAdapter(self.api)
|
| 2663 |
self.binance_futures_with_fallback = FallbackMarketAdapter(
|
| 2664 |
self.binance_futures_adapter,
|
|
@@ -2702,7 +2740,7 @@ class CryptoMonitor:
|
|
| 2702 |
return self.alpha_with_futures
|
| 2703 |
if coin_cfg.get("data_source") == "okx":
|
| 2704 |
return self.okx_adapter
|
| 2705 |
-
return self.
|
| 2706 |
|
| 2707 |
def _futures_symbols(self):
|
| 2708 |
# OKX 标的不走 Binance 强平 websocket;Alpha 如显式开启合约,则订阅其 Binance 合约 symbol。
|
|
@@ -2711,9 +2749,7 @@ class CryptoMonitor:
|
|
| 2711 |
if not coin.get("enable_futures") or not coin.get("futures_symbol"):
|
| 2712 |
continue
|
| 2713 |
source = coin.get("data_source", "binance")
|
| 2714 |
-
if source in ("binance_futures", "binance_alpha")
|
| 2715 |
-
source == "binance" and FUTURES_ENABLED
|
| 2716 |
-
):
|
| 2717 |
symbols.append(coin["futures_symbol"])
|
| 2718 |
return symbols
|
| 2719 |
|
|
|
|
| 89 |
HYPERLIQUID_BASE = os.getenv("HYPERLIQUID_BASE", "https://api.hyperliquid.xyz")
|
| 90 |
ASTER_FAPI_BASE = os.getenv("ASTER_FAPI_BASE", "https://fapi.asterdex.com")
|
| 91 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 92 |
# 配置持久化文件
|
| 93 |
CURRENT_DIR = os.path.dirname(os.path.abspath(__file__))
|
| 94 |
LOCAL_CONFIG_FILE = os.path.join(CURRENT_DIR, "crypto_config.json")
|
|
|
|
| 114 |
return {
|
| 115 |
"symbol": "BTCUSDT",
|
| 116 |
"futures_symbol": "BTCUSDT",
|
| 117 |
+
# 数据来源:binance(默认)、binance_futures、binance_alpha 或 okx。
|
|
|
|
| 118 |
"data_source": "binance",
|
| 119 |
# 首页核心展示:同一时间只应有一个币种开启;旧配置默认 GRAM。
|
| 120 |
"core_display": False,
|
| 121 |
# 休市属性:股票类标的(如 STRC)休市时价格长时间不变,期间不推送任何提醒
|
| 122 |
"market_close": False,
|
| 123 |
"market_close_static_secs": 600,
|
| 124 |
+
# 合约监控按币种开启;币安合约接口(/fapi)与强平 websocket 在部分区域
|
| 125 |
+
# 可能返回 451。受限时会保留现货监控,合约字段降级为空。
|
|
|
|
| 126 |
"enable_futures": False,
|
| 127 |
# 个人关键价位(0 表示未设置,不监控)
|
| 128 |
"cost_price": 0.0,
|
|
|
|
| 1560 |
|
| 1561 |
def __init__(self, primary, fallback, logger, fallback_symbols=None,
|
| 1562 |
primary_label="主源", fallback_label="备用源",
|
| 1563 |
+
allow_symbol_guess=True, fallback_symbol_resolver=None):
|
| 1564 |
self.primary = primary
|
| 1565 |
self.fallback = fallback
|
| 1566 |
self.logger = logger
|
|
|
|
| 1568 |
self.primary_label = primary_label
|
| 1569 |
self.fallback_label = fallback_label
|
| 1570 |
self.allow_symbol_guess = allow_symbol_guess
|
| 1571 |
+
self.fallback_symbol_resolver = fallback_symbol_resolver
|
| 1572 |
self._fallback_logged = set()
|
| 1573 |
|
| 1574 |
+
def _fallback_symbol(self, method, symbol):
|
| 1575 |
+
if self.fallback_symbol_resolver:
|
| 1576 |
+
return self.fallback_symbol_resolver(method, symbol)
|
| 1577 |
if symbol in self.fallback_symbols:
|
| 1578 |
return self.fallback_symbols[symbol]
|
| 1579 |
if self.allow_symbol_guess and symbol.endswith("USDT") and len(symbol) > 4:
|
|
|
|
| 1588 |
if result:
|
| 1589 |
return result
|
| 1590 |
|
| 1591 |
+
fb_symbol = self._fallback_symbol(method, symbol)
|
| 1592 |
if not fb_symbol:
|
| 1593 |
return None
|
| 1594 |
fallback_fn = getattr(self.fallback, method, None)
|
|
|
|
| 1629 |
return self._call("funding_rate_history", symbol, limit)
|
| 1630 |
|
| 1631 |
|
| 1632 |
+
DERIVATIVE_METHODS = {"premium_index", "open_interest", "funding_rate_history"}
|
| 1633 |
+
|
| 1634 |
+
|
| 1635 |
+
def okx_fallback_symbol(method, symbol):
|
| 1636 |
+
text = str(symbol or "").strip().upper()
|
| 1637 |
+
if not text:
|
| 1638 |
+
return None
|
| 1639 |
+
if text.endswith("USDT") and len(text) > 4:
|
| 1640 |
+
asset = text[:-4]
|
| 1641 |
+
suffix = "-USDT-SWAP" if method in DERIVATIVE_METHODS else "-USDT"
|
| 1642 |
+
return f"{asset}{suffix}"
|
| 1643 |
+
return text
|
| 1644 |
+
|
| 1645 |
+
|
| 1646 |
# --- 7. 实时强平监控(websocket,可选)---
|
| 1647 |
class LiquidationTracker:
|
| 1648 |
"""
|
|
|
|
| 1885 |
out["liq_5m_usdt"] = round(max(avg5 * 0.5, 10000))
|
| 1886 |
out["liq_severe_usdt"] = round(out["liq_5m_usdt"] * 4)
|
| 1887 |
|
| 1888 |
+
# 7) 资金费率:历史费率分位(合约启用时)
|
|
|
|
| 1889 |
if (
|
| 1890 |
coin_cfg.get("enable_futures")
|
| 1891 |
and coin_cfg.get("futures_symbol")
|
|
|
|
| 1892 |
):
|
| 1893 |
fr = api.funding_rate_history(coin_cfg["futures_symbol"], 500)
|
| 1894 |
if fr:
|
|
|
|
| 2029 |
self.last_oi_changes = {}
|
| 2030 |
self.last_liq_5m_usdt = None
|
| 2031 |
self.last_rsi = None
|
| 2032 |
+
self.last_data_error = None
|
| 2033 |
|
| 2034 |
def evaluate(self, cfg):
|
| 2035 |
sym = cfg["symbol"]
|
|
|
|
| 2037 |
|
| 2038 |
ticker = self.api.ticker_24h(sym)
|
| 2039 |
k1 = self.api.klines(sym, "1m", 62)
|
| 2040 |
+
ticker_ok = isinstance(ticker, dict) and bool(ticker)
|
| 2041 |
+
k1_ok = isinstance(k1, list) and bool(k1)
|
| 2042 |
+
data_warnings = []
|
| 2043 |
+
if not ticker_ok:
|
| 2044 |
+
data_warnings.append("24h ticker 无返回")
|
| 2045 |
+
if not k1_ok:
|
| 2046 |
+
data_warnings.append("1m K线无返回")
|
| 2047 |
+
elif len(k1) < 17:
|
| 2048 |
+
data_warnings.append(f"1m K线较少 {len(k1)}/17")
|
| 2049 |
+
|
| 2050 |
+
closes = _floats(k1, 4) if k1_ok else []
|
| 2051 |
+
highs = _floats(k1, 2) if k1_ok else []
|
| 2052 |
+
lows = _floats(k1, 3) if k1_ok else []
|
| 2053 |
+
qvols = _floats(k1, 7) if k1_ok else [] # quoteAssetVolume
|
| 2054 |
+
|
| 2055 |
+
def num(value, default=None):
|
| 2056 |
+
try:
|
| 2057 |
+
if value is None or value == "":
|
| 2058 |
+
return default
|
| 2059 |
+
return float(value)
|
| 2060 |
+
except (TypeError, ValueError):
|
| 2061 |
+
return default
|
| 2062 |
+
|
| 2063 |
+
price = num((ticker or {}).get("lastPrice")) if ticker_ok else None
|
| 2064 |
+
if price is None and closes:
|
| 2065 |
+
price = closes[-1]
|
| 2066 |
+
if price is None:
|
| 2067 |
+
self.last_data_error = f"{sym} " + ";".join(data_warnings or ["无可用价格"])
|
| 2068 |
+
return alerts, None
|
| 2069 |
+
self.last_data_error = None
|
| 2070 |
+
|
| 2071 |
+
high24 = num((ticker or {}).get("highPrice"), 0.0) if ticker_ok else 0.0
|
| 2072 |
+
low24 = num((ticker or {}).get("lowPrice"), 0.0) if ticker_ok else 0.0
|
| 2073 |
+
qvol24 = num((ticker or {}).get("quoteVolume"), 0.0) if ticker_ok else 0.0
|
| 2074 |
volume5m = sum(qvols[-5:]) if len(qvols) >= 5 else 0.0
|
| 2075 |
avg5m_volume = qvol24 / 288.0 if qvol24 > 0 else 0.0
|
| 2076 |
volume_ratio_5m = (volume5m / avg5m_volume) if avg5m_volume > 0 else None
|
|
|
|
| 2087 |
"data_source": cfg.get("data_source", self.data_source),
|
| 2088 |
"enable_futures": bool(cfg.get("enable_futures")),
|
| 2089 |
"price": price,
|
| 2090 |
+
"high24": high24 if high24 > 0 else None,
|
| 2091 |
+
"low24": low24 if low24 > 0 else None,
|
| 2092 |
+
"quoteVolume24h": qvol24 if qvol24 > 0 else None,
|
| 2093 |
"volume5m": volume5m,
|
| 2094 |
"volumeRatio5m": volume_ratio_5m,
|
| 2095 |
"priceChange": price_changes,
|
|
|
|
| 2098 |
"openInterest": self.last_open_interest,
|
| 2099 |
"openInterestChange": copy.deepcopy(self.last_oi_changes),
|
| 2100 |
"liquidation5m": self.last_liq_5m_usdt,
|
| 2101 |
+
"dataWarnings": data_warnings,
|
| 2102 |
"updatedAt": datetime.now(timezone.utc).isoformat(),
|
| 2103 |
}
|
| 2104 |
|
|
|
|
| 2201 |
# 11) RSI
|
| 2202 |
self._check_rsi(cfg, sym, head, add)
|
| 2203 |
|
| 2204 |
+
# 12) 合约类
|
| 2205 |
if (
|
| 2206 |
cfg.get("enable_futures")
|
| 2207 |
and cfg.get("futures_symbol")
|
|
|
|
| 2208 |
):
|
| 2209 |
self._check_futures(cfg, head, add)
|
| 2210 |
|
|
|
|
| 2689 |
self.bitget = BitgetAPI(self.log)
|
| 2690 |
self.okx = OKXAPI(self.log)
|
| 2691 |
self.okx_adapter = OKXMarketAdapter(self.okx, self.log)
|
| 2692 |
+
self.binance_spot_with_fallback = FallbackMarketAdapter(
|
| 2693 |
+
self.api,
|
| 2694 |
+
self.okx_adapter,
|
| 2695 |
+
self.log,
|
| 2696 |
+
primary_label="Binance现货主源",
|
| 2697 |
+
fallback_label="OKX",
|
| 2698 |
+
fallback_symbol_resolver=okx_fallback_symbol,
|
| 2699 |
+
)
|
| 2700 |
self.binance_futures_adapter = BinanceFuturesMarketAdapter(self.api)
|
| 2701 |
self.binance_futures_with_fallback = FallbackMarketAdapter(
|
| 2702 |
self.binance_futures_adapter,
|
|
|
|
| 2740 |
return self.alpha_with_futures
|
| 2741 |
if coin_cfg.get("data_source") == "okx":
|
| 2742 |
return self.okx_adapter
|
| 2743 |
+
return self.binance_spot_with_fallback
|
| 2744 |
|
| 2745 |
def _futures_symbols(self):
|
| 2746 |
# OKX 标的不走 Binance 强平 websocket;Alpha 如显式开启合约,则订阅其 Binance 合约 symbol。
|
|
|
|
| 2749 |
if not coin.get("enable_futures") or not coin.get("futures_symbol"):
|
| 2750 |
continue
|
| 2751 |
source = coin.get("data_source", "binance")
|
| 2752 |
+
if source in ("binance", "binance_futures", "binance_alpha"):
|
|
|
|
|
|
|
| 2753 |
symbols.append(coin["futures_symbol"])
|
| 2754 |
return symbols
|
| 2755 |
|
requirements.txt
CHANGED
|
@@ -1,5 +1,6 @@
|
|
| 1 |
fastapi
|
| 2 |
uvicorn[standard]
|
| 3 |
requests
|
|
|
|
| 4 |
python-dotenv
|
| 5 |
websocket-client
|
|
|
|
| 1 |
fastapi
|
| 2 |
uvicorn[standard]
|
| 3 |
requests
|
| 4 |
+
huggingface_hub
|
| 5 |
python-dotenv
|
| 6 |
websocket-client
|