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