Spaces:
Paused
Paused
| 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 | |