Repository navigation
Expand file tree
/
Copy pathmonitor.py
More file actions
executable file
·232 lines (194 loc) · 7.84 KB
/
Copy pathmonitor.py
File metadata and controls
executable file
·232 lines (194 loc) · 7.84 KB
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
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
#!/home/solana/.starknet/venv/bin/python
"""Starknet validator monitor: balance, attestations, containers, sync. Alerts via Telegram."""
import json
import os
import subprocess
import sys
import time
import urllib.parse
import urllib.request
CONFIG_PATH = "/home/solana/.starknet/monitor_config.json"
STATE_PATH = "/home/solana/.starknet/monitor_state.json"
def load_config():
with open(CONFIG_PATH) as f:
return json.load(f)
def load_state():
if not os.path.exists(STATE_PATH):
return {}
try:
with open(STATE_PATH) as f:
return json.load(f)
except Exception:
return {}
def save_state(state):
tmp = STATE_PATH + ".tmp"
with open(tmp, "w") as f:
json.dump(state, f, indent=2)
os.replace(tmp, STATE_PATH)
def tg_send(cfg, text):
url = f"https://api.telegram.org/bot{cfg['tg_token']}/sendMessage"
data = urllib.parse.urlencode({
"chat_id": cfg["tg_chat_id"],
"text": text,
"parse_mode": "HTML",
"disable_web_page_preview": "true",
}).encode()
try:
with urllib.request.urlopen(url, data=data, timeout=10) as resp:
return resp.status == 200
except Exception as e:
log(f"telegram error: {e}")
return False
def log(msg):
print(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] {msg}", flush=True)
def fetch_metrics(url):
metrics = {}
try:
with urllib.request.urlopen(url, timeout=5) as resp:
for line in resp.read().decode().splitlines():
if line.startswith("#") or not line.strip():
continue
name, _, value = line.partition(" ")
key = name.split("{", 1)[0]
try:
metrics[key] = float(value)
except ValueError:
pass
except Exception as e:
log(f"metrics fetch error: {e}")
return metrics
def rpc_call(url, method, params, timeout=10):
body = json.dumps({"jsonrpc": "2.0", "method": method, "params": params, "id": 1}).encode()
req = urllib.request.Request(url, data=body, headers={"Content-Type": "application/json"})
try:
with urllib.request.urlopen(req, timeout=timeout) as resp:
data = json.loads(resp.read().decode())
return data.get("result"), data.get("error")
except Exception as e:
return None, {"message": str(e)}
def container_status(name):
try:
out = subprocess.check_output(
["docker", "inspect", "--format", "{{.State.Status}}", name],
stderr=subprocess.DEVNULL,
timeout=10,
).decode().strip()
return out
except Exception:
return "missing"
def balance_level(balance, cfg):
if balance < cfg["balance_emerg_strk"]:
return "emerg"
if balance < cfg["balance_crit_strk"]:
return "crit"
if balance < cfg["balance_warn_strk"]:
return "warn"
return "ok"
def days_left(balance, burn_per_day=6.0):
return balance / burn_per_day if burn_per_day > 0 else float("inf")
SEVERITY = {"ok": 0, "warn": 1, "crit": 2, "emerg": 3}
def main():
cfg = load_config()
state = load_state()
name = cfg.get("validator_name", "validator")
alerts = []
metrics = fetch_metrics(cfg["metrics_url"])
# 1. Operational balance
bal = metrics.get("validator_attestation_operational_account_balance_strk")
if bal is not None:
cur_level = balance_level(bal, cfg)
prev_level = state.get("balance_level", "ok")
if SEVERITY[cur_level] > SEVERITY[prev_level]:
icons = {"warn": "🟡", "crit": "🔴", "emerg": "⛔"}
alerts.append(
f"{icons[cur_level]} <b>[{name}]</b> Balance {cur_level.upper()}\n"
f"Operational: <b>{bal:.2f} STRK</b> (~{days_left(bal):.1f} days)\n"
f"Threshold: warn {cfg['balance_warn_strk']} / crit {cfg['balance_crit_strk']} / emerg {cfg['balance_emerg_strk']}"
)
elif prev_level != "ok" and cur_level == "ok":
alerts.append(
f"✅ <b>[{name}]</b> Balance recovered\n"
f"Operational: <b>{bal:.2f} STRK</b> (~{days_left(bal):.1f} days)"
)
state["balance_level"] = cur_level
state["balance_last"] = bal
# 2. Missed epochs
missed = metrics.get("validator_attestation_missed_epochs_count")
if missed is not None:
prev = state.get("missed_count", missed)
if missed > prev:
delta = int(missed - prev)
alerts.append(
f"🚨 <b>[{name}]</b> Missed epoch(s)!\n"
f"New misses: <b>+{delta}</b> (total {int(missed)})"
)
state["missed_count"] = missed
# 3. Failed attestations
failed = metrics.get("validator_attestation_attestation_failure_count")
if failed is not None:
prev = state.get("failed_count", failed)
if failed > prev:
delta = int(failed - prev)
alerts.append(
f"🚨 <b>[{name}]</b> Attestation failure(s)!\n"
f"New failures: <b>+{delta}</b> (total {int(failed)})"
)
state["failed_count"] = failed
# 4. Containers
container_state = state.get("containers", {})
new_container_state = {}
for c in cfg["containers"]:
cur = container_status(c)
new_container_state[c] = cur
prev = container_state.get(c)
if prev is not None and prev != cur:
if cur == "running":
alerts.append(f"✅ <b>[{name}]</b> Container <code>{c}</code> back to running")
else:
alerts.append(f"⚠️ <b>[{name}]</b> Container <code>{c}</code> is <b>{cur}</b> (was {prev})")
elif prev is None and cur != "running":
alerts.append(f"⚠️ <b>[{name}]</b> Container <code>{c}</code> is <b>{cur}</b>")
state["containers"] = new_container_state
# 5. Pathfinder sync lag (compare to public RPC)
local_block, _ = rpc_call(cfg["rpc_url"], "starknet_blockNumber", [])
public_block, _ = rpc_call(cfg["public_rpc_url"], "starknet_blockNumber", [])
if isinstance(local_block, int) and isinstance(public_block, int):
lag = public_block - local_block
was_lagging = state.get("lag_alerted", False)
if lag > cfg["block_lag_threshold"] and not was_lagging:
alerts.append(
f"⏰ <b>[{name}]</b> Pathfinder lagging\n"
f"Local: {local_block}, public: {public_block} (diff <b>{lag}</b> blocks)"
)
state["lag_alerted"] = True
elif lag <= cfg["block_lag_threshold"] and was_lagging:
alerts.append(f"✅ <b>[{name}]</b> Pathfinder back in sync (diff {lag})")
state["lag_alerted"] = False
state["last_local_block"] = local_block
# 6. Pathfinder stuck (block not advancing)
last_seen_block = state.get("stuck_block")
last_seen_at = state.get("stuck_block_at", 0)
now = time.time()
if isinstance(local_block, int):
if last_seen_block == local_block:
if now - last_seen_at > cfg["stuck_seconds"] and not state.get("stuck_alerted"):
alerts.append(
f"⏸️ <b>[{name}]</b> Pathfinder stuck\n"
f"Block <code>{local_block}</code> hasn't advanced for {int(now - last_seen_at)}s"
)
state["stuck_alerted"] = True
else:
if state.get("stuck_alerted"):
alerts.append(f"✅ <b>[{name}]</b> Pathfinder advancing again (block {local_block})")
state["stuck_block"] = local_block
state["stuck_block_at"] = now
state["stuck_alerted"] = False
# Send all alerts
for a in alerts:
tg_send(cfg, a)
log(f"alert sent: {a.splitlines()[0]}")
save_state(state)
if not alerts:
log("ok")
if __name__ == "__main__":
main()