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