coinpush / app /core /platform.py
zt p
Add core asset candlestick chart
a4999bf
Raw
History Blame Contribute Delete
5.56 kB
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