Spaces:
Paused
Paused
File size: 5,557 Bytes
3be8894 5dbdea7 3be8894 5dbdea7 ce673a5 5dbdea7 3be8894 5dbdea7 3be8894 5dbdea7 3be8894 5dbdea7 ce673a5 5dbdea7 3be8894 5dbdea7 3be8894 5dbdea7 3be8894 5dbdea7 7dc58b6 5dbdea7 ce673a5 5dbdea7 3be8894 5dbdea7 3be8894 5dbdea7 ce673a5 a4999bf ce673a5 3be8894 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 | import json
import os
import threading
from app.core.database import MonitorStore
from app.core.state import HubStateStore
from app.monitors.crypto_price import CryptoPriceMonitor
from app.monitors.market_move import MarketMoveMonitor
from app.monitors.strategy_treasury import StrategyTreasuryMonitor
class MonitoringPlatform:
"""Runtime container for monitor plugins and shared services."""
def __init__(self):
data_dir = os.getenv("MONITOR_DATA_DIR", "/data")
if not os.path.isdir(data_dir) or not os.access(data_dir, os.W_OK):
data_dir = os.path.join(os.getcwd(), "data")
self.data_dir = data_dir
self.state_store = HubStateStore(data_dir)
self.store = MonitorStore(os.getenv("MONITOR_DB_PATH", os.path.join(data_dir, "monitor.db")))
crypto_config = self.state_store.get("crypto_config")
if crypto_config is not None:
self._write_crypto_config(crypto_config)
self.store.set_kv("crypto_config", crypto_config)
market_move_config = self.state_store.get("market_move:config")
if market_move_config is not None:
self.store.set_kv("market_move:config", market_move_config)
market_move_alert_state = self.state_store.get("market_move:alert_state")
if market_move_alert_state is not None:
self.store.set_kv("market_move:alert_state", market_move_alert_state)
self._lock = threading.RLock()
self.crypto = CryptoPriceMonitor(self.store, self.state_store)
self.strategy_treasury = StrategyTreasuryMonitor(self.store, self.state_store)
self.market_move = MarketMoveMonitor(self.store, self.state_store)
self.plugins = {
"crypto_price": self.crypto,
"strategy_treasury": self.strategy_treasury,
"market_move": self.market_move,
}
paused_ids = self.state_store.get("paused_monitor_ids") or []
if paused_ids:
self.store.apply_paused_ids(paused_ids)
self.started = False
def _write_crypto_config(self, cfg):
try:
from coinpush import CONFIG_FILE
os.makedirs(os.path.dirname(CONFIG_FILE), exist_ok=True)
temp_path = CONFIG_FILE + ".tmp"
with open(temp_path, "w", encoding="utf-8") as handle:
json.dump(cfg, handle, ensure_ascii=False, indent=2)
os.replace(temp_path, CONFIG_FILE)
except Exception as exc:
self.state_store.last_error = str(exc)
def start(self):
with self._lock:
if self.started:
return
for plugin in self.plugins.values():
plugin.start()
self.started = True
def stop(self):
with self._lock:
for plugin in self.plugins.values():
plugin.stop()
self.started = False
self.store.close()
def status(self):
return {
"started": self.started,
"data_dir": self.data_dir,
"state_sync": {
"enabled": self.state_store.enabled,
"remote_synced": self.state_store.remote_synced,
"last_error": self.state_store.last_error,
},
"monitors": self.store.list_monitors(),
"alerts": self.store.recent_alerts(20),
"events": self.store.recent_events(40),
"crypto": {
"status_text": self.crypto.status_text,
"next_wakeup": self.crypto.next_wakeup.isoformat() if self.crypto.next_wakeup else None,
"stats": self.crypto.stats,
"snapshots": self.crypto.snapshots,
"core_market": self.crypto.core_market_summary(),
},
"strategy_treasury": {
"enabled": self.strategy_treasury.enabled,
"interval": self.strategy_treasury.interval,
"stats": self.strategy_treasury.stats,
},
"market_move": {
"enabled": self.market_move.enabled,
"interval": self.market_move.interval,
"threshold_pct": self.market_move.threshold_pct,
"stats": self.market_move.stats,
"snapshots": self.market_move.snapshots,
},
}
def set_monitor_paused(self, monitor_id, paused):
result = self.store.set_monitor_paused(monitor_id, paused)
paused_ids = set(self.state_store.get("paused_monitor_ids") or [])
if paused:
paused_ids.add(monitor_id)
else:
paused_ids.discard(monitor_id)
self.state_store.update(paused_monitor_ids=sorted(paused_ids))
return result
def get_config(self):
return self.crypto.get_config()
def update_config(self, cfg):
result = self.crypto.update_config(cfg)
self.state_store.set_kv("crypto_config", result)
return result
def calibrate_all(self):
self.crypto.calibrate_all()
return self.crypto.get_config()
def calibrate_coin(self, name):
return self.crypto.calibrate_coin(name)
def get_klines(self, coin_name, interval="1h", limit=120):
return self.crypto.get_klines(coin_name, interval, limit)
def get_market_move_config(self):
return self.market_move.get_config()
def update_market_move_config(self, cfg):
result = self.market_move.update_config(cfg)
self.state_store.set_kv("market_move:config", result)
return result
|